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:
- Call
open(consumer_id, wal_dir, state_dir)to open or create the consumer. Keepconsumer_idstable across restarts, give each replica a distinct ID, and putstate_diron durable storage so its cursor survives restarts. - Call
poll()to list uncommitted segments oldest first. It is a directory listing, so it is cheap to call in a loop. - Call
read(&segment)to get CRC-verified ArrowRecordBatches. It finds a segment again if the compactor renames it while the consumer is reading. - 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:
--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, orIcebergSubscription::new, orresumewith no snapshot id). The next poll is a full scan filtered totime_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/sqlover 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_LASTabove the number of commits in your longest expected outage. Compaction and reclustering commits count. Alert onsiglake_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 asincrease(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.