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.