LaserData Cloud
Laser SDK

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:

OptionDefaultWhat it does
batch_length(n)1000messages per server batch
linger(d)0Maximum wait before an incomplete batch is flushed
retries(n, interval)3 retries. Rust waits 250 ms, Python and TypeScript wait 1 sresend attempts on failure, see the notes below for None
routing(r)Balanceddefault routing for every send: Balanced, Routing::key(k), or a fixed partition
create_stream(bool) / create_topic(bool)truecreate missing stream/topic on init, turn off to fail fast against a typo
partitions(n)1partition count used when the producer creates the topic
replication_factor(n)server defaultreplica count when creating the topic
expire_after(d) / never_expire()server defaultmessage retention when creating the topic
max_topic_bytes(n) / unlimited_topic_size()server defaultsize 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 mode None passes through to Iggy, which retries forever. Python and TypeScript default to 3 retries at 1 second.
  • Rust background mode uses BackgroundConfig for batching, shards, and backpressure. OrderedSharding preserves order, while BalancedSharding favors throughput. Call shutdown() to flush buffered messages.
  • Rust's topic.batching() provides explicit flush() and close() without a background task. Defaults are 512 records, 1 MiB, and 5 ms.
  • Python exposes direct-mode fields through batch_length, linger_ms, retries, and retry_interval_ms, plus creation, partition, expiry, and size configuration. Call await producer.init() before sending. Python has no background mode.
  • TypeScript's producer accepts routing, retries, and retryIntervalMs. It has no producer-side batchLength or linger. Use publishBatch() or sendBatch(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 accepts First, Last, Next, Offset(n), or TimestampMicros(t). TypeScript uses { kind: "first" | "last" | "next" | "offset" | "timestamp" }. Python uses polling="first" | "last" | "next" with offset= or timestamp_micros=.
  • allow_replay() permits reads at or below a stored offset. Use it to replay a group from First. Without it, the consumer skips previously consumed records.
  • commit(message) acknowledges a processed record. Native consumers also expose store_offset(offset, partition) and delete_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) and last_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:

PolicyStores the offset
Polling (default)Native consumers: the polled batch end before delivery. Group-policy consumers: the handled prefix at the next poll or shutdown
Allafter consuming everything a poll returned
Eachafter 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
Disablednever, 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 CommitPolicy variants, polling_retry_interval, and init_retries. next_within(timeout) returns a typed timeout error without an outer tokio::time::timeout.
  • Python selects policies through auto_commit. commit_interval_ms adds the interval variant, and commit_every sets the count for every. Call await consumer.init() before iteration.
  • TypeScript's autoCommit: true stores offsets on each poll. With false, use consumer.commit(message). Interval, each, and every policies are unavailable. Defaults are batchLength 100 and pollIntervalMs 250.

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

VerbWhat 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.

ConfigurationEnvironment variableDefaultRange
Report intervalLASER_PRODUCER_TELEMETRY_INTERVAL_MS100000 disables, otherwise 1000 through 60000
Reported handles per connectionLASER_PRODUCER_TELEMETRY_MAX_PRODUCERS320 through 32
CounterMeaning
Submitted records and bytesEach application send once, before retries. Payload bytes only.
Confirmed records and bytesFully successful synchronous send calls. Unavailable for background enqueues.
Failed callsSend calls that returned an error.
Latency p50, p99, p99.9Publish-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.

On this page