Skip to content

Consuming WAL segments

Use siglake_wal::consumer::SegmentConsumer in an external process to read accepted events from sealed write-ahead log (WAL) segments. It keeps a durable cursor and delivers at least once. Its watermark holds retention open within a bounded window. If read reports that a segment is gone, query the missed interval for rows the table still retains, then resume from the current poll.

Do not read the WAL directory directly. The source repository's CONSUMING_SEGMENTS.md is the complete interface and deployment contract.

Choose the right stream

Interface Reads Visibility Delivery model If you fall too far behind
SegmentConsumer Sealed WAL segments Seal cadence, normally seconds Durable cursor, at-least-once; retention waits within a bounded window read fails with "segment is no longer present": the sweep passed your watermark. Backfill the interval by query.
siglake subscribe An Iceberg table Commit cadence Replayable rows from append commits, late arrivals included; non-appends are skipped; refuses when snapshot expiry has broken the chain poll refuses with a typed HistoryGap error and neither half of the cursor moves. Re-bootstrap or backfill.
GET /api/v1/stream Events on one ingester, teed before the WAL append Live ingest path Best-effort push; no history or durability, and slow subscribers are dropped Nothing to recover from the stream. It keeps no history, so a dropped subscriber has no cursor to resume.

Use SegmentConsumer when an external pipeline needs events before they reach Iceberg, in the form Siglake accepted them. Use siglake subscribe when the pipeline should consume a queryable, committed table. Read Tailing a table with siglake subscribe first, because a table subscription does not hold retention open the way a WAL consumer does. Use the SSE stream for a live display that can tolerate gaps and unacknowledged rows. It sends events before the WAL append, so a batch then refused for backlog (503) or whose append fails (500) has still been published. The caller's retry publishes it again. See Ingest and the WAL.

Run the four-call loop

Before you start, choose a stable consumer ID and a durable directory for its cursor. Give each replica a different ID.

use siglake_wal::consumer::SegmentConsumer;

let mut consumer = SegmentConsumer::open(
    "my-consumer",
    "/siglake/wal/default",
    "/var/lib/my-consumer",
)?;

loop {
    for segment in consumer.poll()? {
        for batch in consumer.read(&segment)? {
            process(batch)?;
        }
        consumer.commit(&segment)?; // only after processing succeeds
    }
    std::thread::sleep(std::time::Duration::from_secs(1));
}

The calls form one ordered contract:

  1. Call open(consumer_id, wal_dir, state_dir) to open or create the consumer. Keep consumer_id stable across restarts, give each replica a distinct ID, and put state_dir on durable storage so its cursor survives restarts.
  2. Call poll() to list uncommitted segments oldest first. It is a directory listing, so it is cheap to call in a loop.
  3. Call read(&segment) to get CRC-verified Arrow RecordBatches. It finds a segment again if the compactor renames it while the consumer is reading.
  4. Call commit(&segment) to advance the durable cursor and publish the retention watermark. Call it only after all work for the segment is durable.

Delivery and retention

Delivery is at least once. If the process fails after doing its work but before commit, the segment is delivered again after restart. Make side effects idempotent when duplicates matter. SegmentRef::name is the stable, unique token for deduplication.

A successful commit also tells the compactor how far this consumer has progressed. With the current defaults, committed segments have a 60-second retention floor, a fresh consumer can hold them for up to one hour, and a watermark that has not advanced for five minutes is treated as stale. These bounds keep a stopped consumer from filling the WAL volume.

Alert on consumer lag. Once a consumer exceeds the retention ceiling, the guarantee becomes an explicit "segment is no longer present" read error rather than indefinite disk growth.

One SegmentConsumer reads one segment directory. Run one for each tenant WAL directory. For user indexes, run one for each <wal>/<tenant>/<index>/ directory. Your consumer must divide work across replicas because the WAL does not partition records for external readers.

When the segments you need are gone

Recover available rows from the table, not from the WAL. The sweep only deletes segments the drain has committed, but table retention can later remove their rows. Run POST /api/v1/sql over the interval you missed, then carry on from the current poll. If the table no longer retains those rows, the deleted WAL segments cannot replay them.

The failure arrives from read, not from poll. poll lists what is on disk, and the segment it hands you then fails to open, with an error naming the segment and saying it was swept before this consumer read it. Commit the segments you did process, and read the error as the lag alert you did not act on.

No setting extends the hour above. That ceiling is what keeps one stuck consumer from filling the WAL volume, so the remedy is consumer throughput or replicas, not a longer window.

A table subscription fails differently in kind. Its poll refuses outright with a typed HistoryGap error and neither half of the cursor moves. See When history is gone, the poll refuses.

Tailing a table with siglake subscribe

siglake subscribe (and the IcebergSubscription behind it) tails committed rows from an Iceberg table rather than the WAL:

siglake subscribe --table events --since 2026-09-10T00:00:00Z --interval-secs 5

--table takes events or an index id. --since is where the cursor starts, RFC3339, defaulting to one minute ago; --interval-secs is the poll cadence, default 1; --once processes the current snapshot and exits; --time-column overrides the cursor column, which defaults to events.timestamp or an index's declared timestamp_field. The CLI reference has the rest, including the warehouse and catalog flags.

Its cursor is a pair: the last event time seen and the Iceberg current_snapshot_id last acknowledged. Persist both or neither: a time cursor alone re-bootstraps with a full scan and drops late-arriving rows older than it. The siglake subscribe CLI keeps the pair in memory only, so a restart is always a re-bootstrap from --since.

Each poll walks parent_snapshot_id links back from the table's current snapshot to the cursor's snapshot and reads only the data files those commits added. That backwards walk finds late arrivals committed after your cursor but timestamped before it.

Compaction does not re-deliver rows

A poll interval routinely contains commits that add data files full of rows you already have: the re-clustering compactor merging small files, retention dropping old ones, a delete task rewriting a file without the deleted rows. Every one of those is an Iceberg overwrite commit, and its replacement files are ADDED in the manifest exactly like an append's are. Nothing in the manifest distinguishes a replacement file from a new-row file.

So the subscription classifies each in-range commit by its snapshot summary before reading anything. append commits deliver their files. Every non-append commit is skipped and the cursor advances past it. Compaction with no ingest behind it therefore yields zero rows, and an interval holding both an append and a re-cluster yields exactly the appended rows, including late arrivals.

Each skip increments siglake_subscription_rewrite_commits_skipped_total{table,origin}:

origin Meaning
siglake The commit carries Siglake's siglake.rewrite snapshot-summary property, so Siglake wrote it during compaction, retention, or a delete task. These commits add no rows that the subscription has not already offered. Expect a steady rate on an active table.
foreign The commit is a non-append that Siglake did not stamp. Assume rows were lost and follow the recovery steps below.

External overwrite semantics are not supported

If another engine commits a non-append to a Siglake table, such as a Spark INSERT OVERWRITE, a MERGE, or a row-level delete, the subscription skips it exactly like a compaction commit, counts it with origin="foreign", and logs a WARN naming the table and snapshot ID. Any rows that commit added are not delivered, and the cursor moves past it, so nothing will re-offer them. For anyone dual-writing a Siglake table, that is data loss the consumer cannot see: the only signal is the WARN and the counter on the subscriber, never an error the poll returns.

Recover the same two ways as a history gap: backfill by query with POST /api/v1/sql over the interval the foreign commit covered, then resume from the table's current snapshot ID. You can instead stop dual-writing and write through Siglake's ingest path, whose commits are appends and are delivered normally.

Alert on the series if any writer other than Siglake touches the table, for example increase(siglake_subscription_rewrite_commits_skipped_total{origin="foreign"}[10m]) > 0. One caveat before you treat every foreign skip as lost rows: re-clusters Siglake committed before the siglake.rewrite marker existed land in the same bucket, and those lost nothing. Check the WARN's snapshot ID against the table's commit history before concluding rows are missing.

When history is gone, the poll refuses

The compactor drops old snapshots from table metadata on a timer: SIGLAKE_SNAPSHOT_RETAIN_LAST (default 100), swept every 60 seconds. If a consumer is down for more than retain_last commits, the ancestors between its cursor and the current snapshot can be gone. Those commits cannot be enumerated, so poll() refuses. It returns a typed HistoryGap error (downcastable from anyhow; siglake subscribe exits with it) and advances neither half of the cursor. It does not deliver the reachable suffix because that would move the cursor past commits the consumer never saw. This prevents routine retention from causing silent data loss.

The cursor's own snapshot expiring is not a gap. A retained child whose parent_snapshot_id names the cursor snapshot proves the chain is intact, and delivery continues normally. Only a hole between the two ends stops the subscription.

Recovering from a gap

Table metadata cannot recover the missed interval. Choose one of these actions. The subscription keeps failing until you act:

  • Re-bootstrap. Start a fresh subscription (siglake subscribe --since, or IcebergSubscription::new, or resume with no snapshot id). The next poll is a full scan filtered to time_column > cursor, so you get the missed rows whose timestamps are newer than the cursor. You permanently miss any late-arriving row inside the gap that was timestamped before it.
  • Backfill by query. POST /api/v1/sql over the interval you missed is the complete answer, and the only one that recovers late-arriving rows. Then resume from the table's current snapshot ID.
  • Prevent it. Keep SIGLAKE_SNAPSHOT_RETAIN_LAST above the number of commits in your longest expected outage. Compaction and reclustering commits count. Alert on siglake_subscription_history_gap_total, which increments once per refused poll, per table. The chart has no alert for it, so add your own rule, such as increase(siglake_subscription_history_gap_total[10m]) > 0.

The WAL consumer has the stronger property here: its watermark holds retention open within the bounds above, while a table subscription publishes nothing that snapshot expiry consults. If losing an interval is unacceptable, consume the WAL.

Reference consumer

Siglake's former in-tree semantic detection pipeline, including streaming detectors, episode correlation, and webhook dispatch, moved out of the storage engine and now consumes the WAL entirely through these four calls. It is maintained as the reference consumer: if that pipeline needs a guarantee the interface does not provide, the interface is the part that must change.