Skip to content

Ingest and the WAL

The write-ahead log (WAL) accepts data before it reaches Iceberg. It provides a local durability boundary and a stream of immutable Arrow segments for the compactor and other consumers.

Segment lifecycle

A segment moves through four directories:

active/      open Arrow IPC file
    |
    | seal by event count, age, or explicit request
    v
sealed/      immutable framed segment
    |
    | claimed by a drain
    v
processing/  segment being committed
    |
    | Iceberg commit succeeds
    v
committed/   retained for registered consumers, then swept

Local transitions use filesystem renames. The seal path syncs the segment before renaming it from active/ to sealed/. The current frame format is v2. Its 52-byte header contains the LWAL magic, a format version, a CRC, the minimum and maximum event timestamps, and the owning Iceberg table's UUID. Readers use the timestamps for pruning and reject a bad CRC.

The writer puts the UUID in the header before the segment's first append. A segment recovered from active/ or the object-store mirror therefore keeps its identity. Frame v1 and unframed Arrow IPC segments remain readable, but neither format identifies a table. Ownership checks treat them as having no opinion.

The default segment limits are 4,096 events and five seconds. The --wal-max-events and --wal-max-age-secs flags change them.

Durability model

Ingest routes accept a commit query parameter or X-Siglake-Commit header. The value selects when a successful response is sent.

What is durable when the ingester acknowledges an OTLP export

What has actually been made durable at that point depends on the mode you asked for:

Mode Successful response means
auto The WAL write completed. The bytes can still be only in the kernel page cache.
wait_for (default) The active WAL file and its directory entry were synced before acknowledgement.
force The segment was sealed and the server observed it leave the local sealed/ or processing/ directory before the timeout.

auto opts into an earlier acknowledgement after the WAL write. A node or power failure can lose bytes that remain in the kernel page cache. A successful wait_for acknowledgement means the rows are fsynced on the local WAL and stored nowhere else. They get a copy in an object-store warehouse at whichever comes first: the Iceberg commit, which writes their Parquet files before it publishes the snapshot, or the mirror upload of their segment. The acknowledgement waits for neither boundary. The commit is also what makes the rows visible to queries, unless the query pods mount the WAL buffer.

The WAL sync covers segment bytes and directory entries for the tenant, index, active segment, and sealed segment. Siglake syncs the sealed name before it unlinks the active copy. The power-loss guarantee assumes ext4 or xfs on a node-attached volume. A network filesystem provides whatever its fsync(2) and rename semantics guarantee.

force is unavailable when a catalog-claim drain consumes the mirrored WAL. That remote drain does not move the ingester's local segment, so the ingester cannot observe its commit. The default single-replica local drain supports force while mirroring is on. With a catalog-claim drain, use wait_for and query for visibility. The HTTP API reference lists its timeout and error responses.

Mirroring and disaster recovery

When you configure a warehouse URL, Siglake copies sealed segments to <warehouse bucket>/wal-mirror/ by default. --wal-mirror-prefix selects another prefix, and an empty value disables the mirror. Setting --wal-active-mirror-interval-secs above its default of 0 also copies the current active segment on a timer.

Before it queues a sealed segment, Siglake attempts to create a durable hard link under the WAL's mirror-pending/ directory. A successful link keeps the local bytes after compaction and retention remove the segment's other names. Siglake removes the link after it confirms the remote object. If it cannot create the link, it increments siglake_wal_mirror_failures_total{reason="pin"} and uploads without this protection.

Each ingester checks sealed/ and mirror-pending/ at startup and every 300 seconds. The check resumes uploads after an object-store outage or process exit. Set SIGLAKE_WAL_MIRROR_SWEEP_SECS to change the interval, or to 0 to disable the check. A prolonged upload failure retains pinned bytes and can exhaust the local WAL volume.

Acknowledgements do not wait for either upload. A failed or delayed upload therefore extends the amount of acknowledged but uncommitted data exposed to loss of the WAL volume. Rows the compactor has already committed keep their copy in the warehouse whatever the mirror is doing, so the mirror only ever covers the pre-commit window. Monitor the mirror failure metrics.

siglake wal-recover restores mirrored objects beneath the local WAL root. It recreates the tenant and index directory layout and skips a segment already on disk. See Disaster recovery.

Multi-pod coordination

Catalog-claim mode stores segment ownership in the SQL wal_segments table. try_claim uses an atomic SQL update, so multiple drain processes can compete for work without treating an object-store rename as atomic.

This mode is off by default in the Helm chart. When enabled, ingesters register mirrored segments and drain processes claim them. Scaling describes the required settings.

Multiple ingesters that use the filesystem WAL need a shared ReadWriteMany volume. Each writer uses its own segment name, and tenant routing keeps its WAL under the tenant directory.

The drain

The compactor's drain reads sealed segments, converts their Arrow batches to the destination table schema, writes Parquet, and commits the new files to Iceberg. It marks the source segments committed only after the table commit succeeds.

The drain can keep several commits in flight. Commit-accumulation batching can hold small segments until the 32 MiB target or 10-second age limit is reached. These defaults reduce catalog commits while adding delay before committed storage. Compaction covers that tradeoff.

Backpressure

Ingest refuses work at independent limits instead of holding every request open. Tune rejected ingest to identify the active limit from its response and metrics.

Graceful shutdown

On SIGTERM, the ingest server drains its writer lanes and calls TenantWalRouter::seal_all. This turns buffered active data into sealed segments before exit when the process has enough shutdown time.

The pod's terminationGracePeriodSeconds must cover that work. See Scaling before changing the shutdown settings.

Throughput tuning

The ingest rejection procedure connects each throughput setting to its refusal signal and metric.

Consumer coordination

SegmentConsumer reads sealed segments oldest first. Each consumer has a stable ID, a durable cursor, and its own retention watermark.

The consumer commits its cursor after processing a segment. A crash after processing but before that commit causes re-delivery, so the contract is at least once. Consumers must make their processing idempotent or de-duplicate events.

Committed segments remain available until consumer watermarks and the WAL retention bounds permit a sweep. A stale consumer cannot retain them without a limit. The WAL consumer guide lists the poll, read, commit, and lag behavior.