Skip to content

Consumer

Consumer reads records from feed via cursor (feed-scoped) or stream sequence (stream-scoped). Feed-scoped reads walk sealed commit history. Stream-scoped reads target one stream directly, sealed or open. See Architecture.

StreamStack consumer read path

Read Scopes

Scope API What It Reads
Feed-scoped read(from, mode) / subscription.poll(mode) Sealed commits after cursor
Stream-scoped read(feed, streamSeq, mode) One stream by sequence, sealed or unsealed

Feed-scoped reads are continuous-consumption. Cursor marks position in commit history; each page returns records from commits after that position plus end cursor for next call.

Stream-scoped reads bypass commit history. Sealed stream served from commit artifacts. Unsealed stream stitched from durable WAL tail. Inspect open ingest or replay one stream in isolation.

Cursor & Paging

CursorToken is feed name plus commit sequence. Sequence 0 is feed start. from cursor is exclusive: records begin in first commit after that sequence.

API Purpose
head(feed) Current feed head cursor
log(feed) Full sealed commit list
read(from, mode) One page from from to head
read(from, to, mode) One page from from to to

Page shape is ReadPage:

Field Meaning
records Page records
end Cursor after last included commit

Pass end as next from to continue. When end equals head, consumer is caught up.

Paging is bounded server-side:

Setting Value
Default limit 1,000 records
Maximum limit 10,000 records

For RAW and DEDUP, server fills page commit by commit. For LATEST, server resolves winners across full requested range before returning.

Subscriptions

subscription is named, persisted cursor in metadata. Multiple consumers use different names on same feed and advance independently.

Step API What Happens
Bind PUT /v1/subscriptions/{name}?feed=... Associates name with feed; creates cursor at start if new
Poll poll(mode) Reads from saved position to current head
Advance advance(end) Persists new position

Java client binds lazily on first position() or poll() call. Subscription with no saved position starts at commit 0.

Poll returns same ReadPage shape as direct feed read. After processing records, call advance(batch.end()) so next poll skips handled data. Position survives process restarts and re-binding under same name.

Long Poll

poll(mode, wait) blocks up to wait for new commit on feed. Server wakes waiters when seal publishes commit. If nothing arrives before timeout, poll returns empty page at current position.

Long poll avoids tight loops when waiting for next sealed commit.

Read Modes

Modes apply at read time over chosen scope and range. Stored data is unchanged.

Mode Behavior
RAW Every record in arrival order, including duplicates
DEDUP First occurrence of each fingerprint in range
LATEST One winner per key by version (then commit order); keyless records pass through

DEDUP keys on stored fingerprint bytes, not business key. Default fingerprint is SHA-256 of the payload. Identical payloads collapse unless the producer supplied different fingerprints.

LATEST requires keyed feeds for per-key collapse. Uses key index and version resolver: higher version wins, ties break toward later commit. Tombstone winners emitted as tombstone records. Point lookup of tombstoned key returns empty.

Point Lookup

Keyed feeds: latest(feed, key) returns current winning record for one key without scanning range. Server walks key-index metadata and fetches winning pack entry only.

Diff & Merge

Commit-scoped utilities compare keyed state between two cursor positions:

API Purpose
diff(from, to) Added, updated, and deleted keys between two commits
merge(left, right) Three-way merge resolution between two commit tips

Operates on key metadata only; full record payloads not required.

Truncation

Retention and compaction remove old commits from active log. Cursor before log floor fails with CommitTruncatedException.

Situation Result
Cursor precedes retention floor Read rejected; log starts at higher commit
Stream sequence compacted away Stream read rejected; active log starts at later stream

Advance subscriptions and cursors as data is processed so positions stay inside retained window.

Examples

CLI feed-scoped read from start:

./streamstack read orders --from 0 --mode LATEST

CLI stream-scoped read of one open or sealed stream:

./streamstack read stream orders 3 --mode RAW

CLI subscription poll with auto-advance:

./streamstack subscribe etl --feed orders --mode RAW --wait 5000

Java SDK continuous consumption:

Subscription sub = client.subscription("etl", "orders");
ReadPage page = sub.poll(ReadMode.RAW, Duration.ofSeconds(5));
for (Record record : page.records()) {
    process(record);
}
sub.advance(page.end());

Java SDK manual cursor paging:

CursorToken cursor = CursorToken.start("orders");
while (true) {
    ReadPage page = client.reads().read(cursor, ReadMode.DEDUP);
    if (page.records().isEmpty()) {
        break;
    }
    page.records().forEach(this::process);
    cursor = page.end();
}

Java SDK stream-scoped read and point lookup:

StreamReadResponse stream = client.reads().read("orders", 3L, ReadMode.RAW);
stream.records().forEach(this::process);

client.reads().latest("orders", "policy-1").ifPresent(this::process);