Log
Publish, route, consume, and replay messages with explicit ordering and offset policies
The log stores messages for live reading, recovery from an offset, or replay from the beginning. An offset identifies a message's position in a partition. Events and commands use the same message foundation.
Built for
Use Log for event-driven services, audit trails, and live feeds.
How it works
A stream groups topics. A topic is a named log without a required SDK schema. Partitions store its messages in order.
One connection reaches every stream. Use laser.stream("shop").topic("orders"), or set a default stream to omit .stream(..). See Connect.
Order applies within a partition. There is no total order across a topic's partitions. Publication supports three routing choices:
- Balanced, the default, distributes messages for throughput. Only messages in the same partition have a shared order.
- Key routing hashes a non-empty key to a partition. Messages with that key retain their relative order.
- Partition routing selects a partition number directly.
More partitions allow more parallel consumption. Fewer partitions keep more messages within the same ordered sequence. Use a customer or conversation key when those messages need a shared order.
Payloads are bytes. The raw operation imposes no encoding. JSON, MessagePack, and Avro helpers encode values into those bytes. These examples use JSON.
A named cursor decodes records and tracks per-partition offsets in its reader handle. Save offsets() and restore with from_offsets(saved) after restart. Without restored offsets, a new reader starts at the beginning. The quick example uses a cursor.
A live consumer continuously reads one partition or joins a consumer group across processes. Each partition has one assigned member. Membership changes cause the server to rebalance assignments. A member that stops reporting is dropped after the session timeout, 30 seconds by default, and its partitions go to the remaining members. Use this path for production processing.
Delivery is at least once. A crash between delivery and offset commit causes redelivery, so handlers must be safe to repeat. Each consumer's commit policy determines when it stores offsets. See Consumers, offsets, and commit policies.
Quick example
const topic = laser.stream("shop").topic("orders")
await topic.ensure(2)
for (const order of ORDERS) {
await topic.publish().json(order).send()
}
const replay = await topic
.json(ORDER_CODEC)
.records("log-example")
for (const result of await replay.poll()) {
if (result.kind === "record") {
console.log(result.record.value.total)
}
}
// One poll reads at most one configured batch per partition, so a real
// reader drains until it reaches the current tail.let topic = laser.stream("shop").topic("orders");
topic.ensure(2).await?;
for order in [Order { id: 1, total: 99 }, Order { id: 2, total: 42 }] {
topic.publish().json(&order)?.send().await?;
}
let mut replay = topic.json::<Order>().records("log-example")?;
while let Some(next) = replay.next().await {
println!("{}", next?.value.total);
}topic = laser.stream("shop").topic("orders", cls=Order)
await topic.ensure(partitions=2)
for order in (Order(id=1, total=99), Order(id=2, total=42)):
await topic.publish(order).send()
reader = topic.records("log-example")
while (record := await reader.next()) is not None:
print(record.value.total)The runnable examples define Order, ORDERS, and TypeScript's ORDER_CODEC. For production, use a tuned producer and a consumer group with a commit policy.
Complete examples are available in Rust, Python, and TypeScript.
Producing at volume
publish() builds one message, sends it, and waits for the result. It uses the connection's publish timeout and retries. See Connect. For higher throughput, select a write path:
publish_batch()collects messages on the client and sends them in one request. You control when the batch is sent.topic.producer()creates a long-lived producer that combines sends into batches and retries failures. Use it as the default production path.- Rust background mode uses Apache Iggy's buffered pipeline across shards for maximum throughput.
The producer builder supports these fields and defaults:
| Option | Default | What it does |
|---|---|---|
batch_length(n) | 1000 | messages per server batch |
linger(d) | 0 | Maximum wait before an incomplete batch is flushed |
retries(n, interval) | 3 retries. Rust waits 250 ms, Python and TypeScript wait 1 s | resend attempts on failure, see the notes below for None |
routing(r) | Balanced | default routing for every send: Balanced, Routing::key(k), or a fixed partition |
create_stream(bool) / create_topic(bool) | true | create missing stream/topic on init, turn off to fail fast against a typo |
partitions(n) | 1 | partition count used when the producer creates the topic |
replication_factor(n) | server default | replica count when creating the topic |
expire_after(d) / never_expire() | server default | message retention when creating the topic |
max_topic_bytes(n) / unlimited_topic_size() | server default | size cap when creating the topic |
This example comes from native-streaming:
await using producer = topic.producer({
retries: 3,
retryIntervalMs: 1_000
})
await producer.send(utf8("message-0"), {
key: utf8("account-42"),
headers: { type: HeaderValue.uint16(7) }
})
await producer.sendBatch(payloads)let producer = topic
.producer()
.batch_length(1_000)
.linger(Duration::from_millis(5))
.retries(Some(3), Some(Duration::from_secs(1)))
.routing(Routing::Balanced)
.build()
.await?;
producer
.send_keyed(
ProducerMessage::new(b"message-0".as_slice())
.header(HeaderKey::try_from("type")?, HeaderValue::from(7_u16)),
b"account-42".to_vec(),
)
.await?;
producer.send_batch(batch).await?;producer = topic.producer(
batch_length=1000,
linger_ms=5,
retries=3,
retry_interval_ms=1000,
)
await producer.init()
await producer.send(
b"message-0",
headers={"type": ("uint16", 7)},
key=b"account-42",
)
await producer.send_batch(values)A send can override the builder's routing. Rust uses send_keyed(msg, key) or send_to_partition(msg, n). Python's send accepts key= or partition=. TypeScript's send accepts { key } or { partition }. Do not supply a key and partition together.
The clients differ in these ways:
- Rust's producer defaults to the connection's publish retry count and delay. In direct mode
retries(None, ..)means no retries. In background modeNonepasses through to Iggy, which retries forever. Python and TypeScript default to 3 retries at 1 second. - Rust background mode uses
BackgroundConfigfor batching, shards, and backpressure.OrderedShardingpreserves order, whileBalancedShardingfavors throughput. Callshutdown()to flush buffered messages. - Rust's
topic.batching()provides explicitflush()andclose()without a background task. Defaults are 512 records, 1 MiB, and 5 ms. - Python exposes direct-mode fields through
batch_length,linger_ms,retries, andretry_interval_ms, plus creation, partition, expiry, and size configuration. Callawait producer.init()before sending. Python has no background mode. - TypeScript's producer accepts
routing,retries, andretryIntervalMs. It has no producer-sidebatchLengthorlinger. UsepublishBatch()orsendBatch(payloads)to combine requests.
Consumers, offsets, and commit policies
Next does not skip the first record of a fresh consumer. With no stored checkpoint, it starts at the first retained record. With a stored checkpoint at offset 0, it starts at offset 1. Writing offset 0 is an acknowledgment of that record, not an empty initial state.
Rust and Python expose last_stored_offset as local progress. On the native path, the Iggy cache can contain an initial zero before a durable commit. On the group-policy path, it reports the last acknowledged offset. Use the server's Next strategy to resume, rather than calculating a start from this diagnostic cache.
Rebuild a native consumer after a purge. A purge restarts every partition at offset 0 and clears every stored offset. An open native consumer can keep its old position and skip replacement records. Shut it down and build a new one with the same name or group. Group-policy readers detect changed source history, invalidate old delivery handles, and reset partition progress. They report source_changed before a subsequent read resumes. A numeric group handle refuses a recreated stream or topic and requires a new reader.
topic.consumer(name, partition) builds a named reader for one partition. topic.consumer_group(group) returns a group handle. Build its reader through group.consumer() to join the group. The server assigns each partition to one member and rebalances as membership changes.
create_group(true) and auto_join_group(true) create and join groups by default. The server stores offsets by consumer name or group. Reconnect with the same name to resume from the stored offset.
Consumer position uses these operations:
start_at(..)selects the initial position when no relevant offset exists. Rust acceptsFirst,Last,Next,Offset(n), orTimestampMicros(t). TypeScript uses{ kind: "first" | "last" | "next" | "offset" | "timestamp" }. Python usespolling="first" | "last" | "next"withoffset=ortimestamp_micros=.allow_replay()permits reads at or below a stored offset. Use it to replay a group fromFirst. Without it, the consumer skips previously consumed records.commit(message)acknowledges a processed record. Native consumers also exposestore_offset(offset, partition)anddelete_offset(partition)for direct offset control. Group-policy consumers refuse those operations because their acknowledgments must pass the group and policy checks.last_consumed_offset(p)andlast_stored_offset(p)read local progress.
To read only some records, use Consumer Filters. The server selects matching records and your progress still covers every record it scanned.
Automatic commit policies control when the consumer stores offsets on the server:
| Policy | Stores the offset |
|---|---|
Polling (default) | Native consumers: the polled batch end before delivery. Group-policy consumers: the handled prefix at the next poll or shutdown |
All | after consuming everything a poll returned |
Each | after every yielded record |
Every(n) | after every n yielded records |
Interval(d) | on a timer |
IntervalOrPolling(d) / IntervalOrAll(d) / IntervalOrEach(d) | the interval or the event, whichever fires first |
Disabled | never, you call commit(message) yourself |
Native polling can skip unprocessed records after a failure because its default policy commits before delivery. Group-policy consumers acknowledge the delivered prefix when consumption continues or shuts down. A crash can cause redelivery. Use Disabled and commit after successful processing when only explicitly completed work must advance. Native consumers send an offset store, while group-policy consumers send a fenced acknowledgment.
await using consumer = await topic.consumerGroup("workers").consumer({
batchLength: 100,
autoCommit: false,
startFrom: { kind: "first" },
pollIntervalMs: 5
})
for (;;) {
const message = await consumer.nextWithin(5_000, { signal })
if (message === null) break
handle(message)
await consumer.commit(message)
}let mut consumer = topic
.consumer_group("workers")
.consumer()
.batch_length(100)
.poll_interval(Duration::from_millis(5))
.start_at(ConsumerStart::First)
.allow_replay()
.commit_policy(CommitPolicy::Disabled)
.build()
.await?;
while let Some(message) = consumer.next().await {
let message = message?;
handle(&message)?;
consumer.commit(&message).await?;
}
consumer.shutdown().await?;consumer = topic.consumer_group("workers").consumer(
batch_length=100,
poll_interval_ms=5,
polling="first",
auto_commit="disabled",
allow_replay=True,
)
await consumer.init()
async for message in consumer:
handle(message)
await consumer.commit(message)
await consumer.shutdown()For the batched default, remove manual commits. Use commit_policy(CommitPolicy::IntervalOrEach(Duration::from_secs(1))) in Rust or auto_commit="each", commit_interval_ms=1000 in Python.
Consumer APIs differ by language:
- Rust provides all ten
CommitPolicyvariants,polling_retry_interval, andinit_retries.next_within(timeout)returns a typed timeout error without an outertokio::time::timeout. - Python selects policies through
auto_commit.commit_interval_msadds the interval variant, andcommit_everysets the count forevery. Callawait consumer.init()before iteration. - TypeScript's
autoCommit: truestores offsets on each poll. Withfalse, useconsumer.commit(message). Interval, each, and every policies are unavailable. Defaults arebatchLength100 andpollIntervalMs250.
shutdown() stops polling and leaves the group. Do not treat it as proof that every automatically committed record was processed. Rust and Python store a delivered offset zero before automatic shutdown when the native SDK reports an absent or zero stored offset. This applies to every automatic commit policy. With Disabled, progress follows explicit commits after processing.
Unnamed TypeScript consumers have isolated identities and default to no automatic commit. Use a stable name and an explicit commit policy to resume durable progress.
Receive only matching records
Keep one shared topic and deliver a different subset to each consumer. Consumer Filters select records on the server, across the topic's partitions, before they cross the network. This reduces payload transfer and decoding work in applications. The complete CDC example saves 98.5% of payload transfer.
Configure the policy through topic.consumer_group(name).filter(). Normal group consumers and advanced group.reader() calls apply it automatically. An unbound group receives all records. Both preserve original payloads and offsets. Raw Iggy accessors retain standard polling behavior and receive every record of a filtered group, so don't mix them with Laser consumers on one group. Typed header policies can select binary payloads without decoding them.
Apache Iggy access
Rust's topic.iggy_producer(), topic.iggy_consumer(..), and topic.iggy_consumer_group(group) return native Iggy builders on the same Iggy connection. Use them for configuration that Laser SDK does not expose.
Key operations
| Verb | What it does |
|---|---|
stream(name).topic(name) | Address a topic on any stream |
topic.ensure(partitions) | Create the topic if it doesn't exist, with a fixed partition count |
publish().payload(bytes) | Append one message, raw and unencoded |
.json(v) / .msgpack(v) / .avro(v) | Encoding convenience over .payload |
.partition_key(key) | Route by key for per-key ordering |
.index(key, value) | Stamp an explicit indexed header, so a view can query it with no projection schema |
.header(key, value) | Attach non-indexed metadata (trace id, source) that rides through to a projector's Row.metadata |
publish_batch() | Accumulate several messages into one round trip |
topic.producer() | A long-lived producer with batching, linger, retries - see Producing at volume |
topic.batching() | A client-side accumulator with explicit flush()/close() (Rust) |
topic.json::<T>().records(name) | A named, resumable cursor |
topic.consumer(..) / topic.consumer_group(name).consumer() | A live stream, partitioned across group members |
consumer.next() / next_within(timeout) | Wait for the next record, optionally bounded by a typed timeout |
consumer.commit(message) | Explicit offset commit for commit-after-handled semantics |
store_offset(..) / delete_offset(..) | Native consumers only. Group-policy consumers use fenced acknowledgments |
topic.replay() | A cursor that reads again from the start instead of tailing |
This is a selection of operations. Views explains the extraction and projection registration behind .index(..).
Durability
Publish completion follows the topic's durability policy. Local queue admission and an HTTP dispatch response do not prove persistence. Message durability and explicit consumer-offset durability are independent. See Server & Durability.
Producer statistics
Producers can report statistics over a separate observer connection. Set the variables before creating producers. A telemetry error never fails a publish.
| Configuration | Environment variable | Default | Range |
|---|---|---|---|
| Report interval | LASER_PRODUCER_TELEMETRY_INTERVAL_MS | 10000 | 0 disables, otherwise 1000 through 60000 |
| Reported handles per connection | LASER_PRODUCER_TELEMETRY_MAX_PRODUCERS | 32 | 0 through 32 |
| Counter | Meaning |
|---|---|
| Submitted records and bytes | Each application send once, before retries. Payload bytes only. |
| Confirmed records and bytes | Fully successful synchronous send calls. Unavailable for background enqueues. |
| Failed calls | Send calls that returned an error. |
| Latency p50, p99, p99.9 | Publish-call time including retry delays, power-of-two bucket upper bounds. p99 needs 100 calls, p99.9 needs 1000. |
- A failed batch can have committed a prefix the confirmed counters leave out. The counters are SDK reports, not delivery totals, exactly-once evidence, or billing input.
- A report expires after three intervals. Handles past the limit publish normally and are left out.
- The Producers tab of a topic in Stream UI shows the connected node only. The observer node can differ from the publishing node. Unknown values show as unavailable, never as zero.
Running it
Log runs on Laser Stack, LaserData Cloud, and standalone Iggy. Context, folded Memory, and Fabric use the same streaming path. Views, Changes, State, and Graph require laser-plane. Consumers must tolerate redelivery.
Related
Consumer Filters
Next: read only the records a filter selects
Views
Queries that already ran
Laser Stack
Choose standalone Iggy or the complete local stack
Quickstart
The same publish/read-back story, end to end
Connect
Connection strings, tokens, and TLS
native-streaming example
Tuned producer, both commit patterns, 1000 messages, all three languages