Messages in depth
Producers, batching, publish failures, consumers, offsets, commit policies, and producer statistics
The full reference for topics and messages. For a short introduction, start with Messages.
Topics, partitions, and routing
A stream groups topics. A topic is a named log with no schema required by the SDK. Partitions store its messages in order. One connection reaches every stream. Address a topic with laser.stream("telemetry").topic("metrics"), or set a default stream and use laser.topic("metrics"). See Connect.
Order holds within a partition. There is no total order across a topic's partitions. A publish picks a partition in one of three ways:
- Balanced, the default, spreads messages for throughput. Only messages in the same partition share an order.
- Key routing hashes a non-empty key to a partition. Messages with the same key keep their relative order.
- Partition routing names a partition number directly.
More partitions allow more parallel consumers. Fewer partitions keep more messages in one ordered sequence. Use a host, agent, or session key when those messages need a shared order.
Payloads are bytes. The raw publish adds no encoding. JSON, CBOR, MessagePack, BSON, and Avro helpers encode values into those bytes. TypeScript snippets on this page turn text into bytes with new TextEncoder().encode(text).
Typed handles and readers
topic.json(..) returns a typed handle. Its publish encodes a value and its records(name) reader decodes each message back. topic.cbor(..) does the same in CBOR. Rust takes a type parameter, topic.json::<Reading>(). Python takes a class, topic.json(Reading). TypeScript takes a codec, topic.json(new Json<Reading>(decode)), whose decode function checks each value because types do not exist at runtime. Cbor, Msgpack, and Bson take the same decode function.
A Python typed publish(body) returns the publish builder, so finish with .send(). A Rust typed publish(&value) also returns a builder. A TypeScript typed publish(value) sends at once and takes { key }, { partition }, or { headers } as a second argument.
A typed reader is a cursor that you own. It is not consumer group delivery, and the server does not store its offsets. It starts at offset 0 unless you restore saved offsets:
| Rust | Python | TypeScript | |
|---|---|---|---|
| Read one record | reader.next().await | await reader.next() | await reader.next() |
| At the current tail | None | None | undefined |
| Read a batch | reader.poll().await | await reader.poll() | await reader.poll() |
| Save offsets | reader.offsets(), a list with one entry per partition | reader.offsets, a list with one entry per partition | reader.offsets, a map from partition to bigint |
| Resume | records(name)?.from_offsets(saved) | records(name, from_offsets=saved) | (await records(name)).fromOffsets(saved) |
| Records per request and partition | 1000 | 1000 | 1000 |
| Record position | MessageId | MessageId with partition_id and offset | { partitionId, offset } |
Change the request size with batch(n) in Rust and TypeScript and records(name, batch=n) in Python. One poll drains each partition to its tail, up to 10,000 records per partition, and returns the records in log timestamp order. The offsets are empty until the first poll or until you restore saved ones. Restoring replaces the offsets, it does not merge them.
A record that does not decode fails with a decode error that names its log position, and the next read moves past it. Rust next() returns the error, TypeScript next() returns a result with kind: "error", and Python next() raises TypedDecodeError. poll() returns it inline next to the good records in all three. A Python async for loop ends on the error, so loop over next() to skip it. Rust and TypeScript stream() and Python async for end once the reader is caught up. topic.replay() is the raw cursor with the same offset rules, across every partition.
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 more throughput, pick a write path:
publish_batch()collects messages on the client and sends them in one request. You decide when the batch goes out.topic.producer()creates a long-lived producer that sends each batch directly and retries failures. Use it as the default production path.- A producer in background mode queues sends and writes them from a worker. It gives the most throughput. See Background producers.
topic.batching()returns a client-side accumulator that flushes on a record count, a byte size, or a timer. See Batching producers.
Direct producer
The Rust producer builder takes these options. Python takes the same values as keyword arguments, and TypeScript takes them as options.
| Option | Default | What it does |
|---|---|---|
batch_length(n) | 1000 | Most messages in one server request. A larger send_batch is split into requests of this size, sent one after another |
linger(d) | 0 | Minimum gap between sends. A send waits out what is left of it since the previous one. Direct mode only |
retries(n, interval) | The connection's publish retries, 3 retries with a first delay of 250 ms | Resend attempts after a failure |
retry_backoff(d) | The connection's first retry delay | Changes the delay and keeps the retry count |
routing(r) | Balanced | Default routing for every send: Routing::Balanced, Routing::key(k), or Routing::Partition(n) |
create_stream(bool) / create_topic(bool) | true | Create a missing stream or topic on init. Turn off to fail fast on a wrong name |
partitions(n) | 1 | Partition count when the producer creates the topic |
expire_after(d) / never_expire() | Server default | Message retention when the producer creates the topic. Setting both forms, or a zero value, fails with an invalid error before any network call |
max_topic_bytes(n) / unlimited_topic_size() | Server default | Size limit when the producer creates the topic. Setting both forms, or a zero value, fails with an invalid error before any network call |
background(config) | Direct mode | Switch to Apache Iggy's buffered mode. batch_length and linger are then ignored |
Python's topic.producer(..) takes batch_length, linger_ms, retries, retry_interval_ms, key or partition for the default routing, create_stream, create_topic, partitions, expire_after_ms or never_expire, max_topic_bytes or unlimited_topic_size, and background. linger_ms and expire_after_ms take a float, so a fraction of a millisecond is kept. max_topic_bytes is a positive byte count. send and send_batch initialize the producer on first use. Call await producer.init() when startup must fail before the service accepts work. send_batch takes raw payloads or (payload, headers) pairs.
TypeScript's topic.producer(options) takes routing, retries, retryBackoffMs, createStream, createTopic, partitions, expireAfterMs or neverExpire, batchLength, lingerMs, maxTopicBytes or unlimitedTopicSize, and background. expireAfterMs is a number of milliseconds and maxTopicBytes is a bigint. The producer needs no init call. It creates a missing stream and topic before the first send unless you turn that off. A transport that cannot set topic expiry or size limits refuses those options.
This sample follows the native-streaming example:
import { HeaderValue } from "@laserdata/laser-sdk"
const text = new TextEncoder()
await using producer = topic.producer({
batchLength: 1_000,
lingerMs: 5,
retries: 3,
retryBackoffMs: 1_000
})
await producer.send(text.encode("reading-0"), {
key: text.encode("node-7"),
headers: { type: HeaderValue.uint16(7) }
})
const payloads = [text.encode("reading-1"), text.encode("reading-2")]
await producer.sendBatch(payloads)use laser_sdk::prelude::*;
use laser_sdk::stream::{HeaderKey, HeaderValue};
use std::time::Duration;
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"reading-0".as_slice())
.header(HeaderKey::try_from("type")?, HeaderValue::from(7_u16)),
b"node-7".to_vec(),
)
.await?;
let batch = vec![
ProducerMessage::new(b"reading-1".as_slice()),
ProducerMessage::new(b"reading-2".as_slice()),
];
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"reading-0",
headers={"type": ("uint16", 7)},
key=b"node-7",
)
values = [b"reading-1", b"reading-2"]
await producer.send_batch(values)A send can override the producer's routing. Rust uses send_keyed(msg, key) or send_to_partition(msg, n). Python's send takes key= or partition=. TypeScript's send takes { key } or { partition }. Do not pass a key and a partition together. Batches take the same override through Rust send_batch_with_routing, Python send_batch(key=, partition=), and TypeScript sendBatch(messages, { key }). All records in one batch share the key or partition.
Producers use the connection's publish retry count and delay unless you set them. In direct mode, retry delays start at the configured delay, 250 ms by default, double after each failure, and stop at 30 seconds. Only temporary errors are retried. In background mode the first resend goes at once, and later resends wait the configured delay each time without doubling. Background mode resends after any error except a lost confirmation, which means the batch may have been written. Each producer setup attempt also uses the publish timeout and retry budget. A producer that loses a race to create the same stream or topic treats it as success and reads the winner's resource. These are the producer-level overrides:
| Goal | Rust | Python | TypeScript |
|---|---|---|---|
| Inherit the connection values | Leave retries unset | retries=None and retry_interval_ms=None, the defaults | Leave retries and retryBackoffMs unset |
| Disable resends | retries(Some(0), None) | retries=0 | retries: 0 |
| Change only the delay | retry_backoff(d) | retry_interval_ms= with retries=None | retryBackoffMs without retries |
In Rust direct mode, retries(None, ..) also means no retries.
Creation settings apply only to topics that do not exist yet. A producer never changes an existing topic's expiry or size limit. A topic created by topic.ensure(..) never expires, and ensure leaves an existing topic as it is. Two producers that create the same topic at once both succeed in all three SDKs, because the one that loses the race accepts the existing topic.
Background producers
Background mode queues each send and writes it from a worker. A send returns once the queue accepts the records, so it carries no commit confirmation, and the publish timeout does not cover that step. Drain the queue before you exit, or the queued records are lost.
All three SDKs take the same settings, with Apache Iggy's background defaults: one ordered shard, a flush at 1000 sends, 1 MiB, or the linger time, a 32 MiB buffer budget, and the block failure mode. Ordered sharding, the default, keeps every send of the producer on one shard in send order. Balanced sharding spreads sends over every shard and gives up that order.
| Setting | Rust BackgroundConfig::builder() | Python BackgroundConfig(..) | TypeScript background: { .. } |
|---|---|---|---|
| Shards | num_shards(n) | num_shards= | shards |
| Sharding | sharding(Box::new(BalancedSharding::default())) | sharding="balanced" | sharding: "balanced" |
| Flush at sends | batch_length(n) | batch_length= | batchLength |
| Flush at bytes | batch_size(n) | batch_size= | batchBytes |
| Linger | linger_time(d) | linger_ms= | lingerMs |
| Buffer budget | max_buffer_size(n) | max_buffer_size= | maxBufferBytes |
| Writes in flight | max_in_flight(n) | max_in_flight= | maxInFlight |
| Full buffer | failure_mode(..) | failure_mode= and block_timeout_ms= | failureMode |
| Failed write | error_callback(..) | error_callback= | onError |
Python selects the mode with background=True for the defaults or background=BackgroundConfig(..) to change them. Rust's builder takes Apache Iggy types: linger_time takes an IggyDuration (convert a Duration with .into()), max_buffer_size takes an IggyByteSize, failure_mode takes a BackpressureMode from laser_sdk::iggy::clients::producer_config, and error_callback takes an ErrorCallback from laser_sdk::iggy::clients::producer_error_callback. Python's linger_ms takes a float, so a fraction of a millisecond is kept.
Close a background producer with shutdown(). It waits for sends in flight and writes every queued record before it closes. After shutdown, a send raises InvalidError in Python and TypeScript. A second shutdown() returns once the first one finishes. In Rust, shutdown() takes the producer by value and needs the last live handle.
By default, Rust and Python log a failed background write and drop its records. Set the error callback to handle it instead. Rust's callback receives Apache Iggy's ErrorCtx with the cause, the stream and topic, the confirmed ranges, and the messages. A Python error_callback receives a PublishFailedError with stream, topic, committed, unconfirmed, and the cause as __cause__. It can be a plain or async callable. TypeScript hands each failure to onError and awaits a returned promise. Without onError, or when it throws, TypeScript keeps the failure and shutdown() rejects with one PublishFailedError that lists the records of every failed background batch.
const text = new TextEncoder()
const producer = topic.producer({
background: {
shards: 2,
lingerMs: 5,
onError: (error) => console.error(error.publishCause())
}
})
await producer.send(text.encode("reading-0"))
await producer.send(text.encode("reading-1"))
await producer.shutdown()use laser_sdk::stream::BackgroundConfig;
use std::time::Duration;
let producer = topic
.producer()
.background(
BackgroundConfig::builder()
.num_shards(2)
.linger_time(Duration::from_millis(5).into())
.build(),
)
.build()
.await?;
producer.send(b"reading-0".as_slice()).await?;
producer.send(b"reading-1".as_slice()).await?;
producer.shutdown().await?;import laser_sdk as ls
def report(error: ls.PublishFailedError) -> None:
print(f"{len(error.unconfirmed)} record(s) unconfirmed: {error.__cause__}")
producer = topic.producer(
background=ls.BackgroundConfig(num_shards=2, linger_ms=5, error_callback=report)
)
await producer.init()
await producer.send(b"reading-0")
await producer.send(b"reading-1")
await producer.shutdown()Batching producers
topic.batching() returns a producer that queues records on the client and sends each flush as one append. A flush runs when the queue reaches max_records, when the queued payload reaches max_bytes, or when linger expires, whichever comes first. A linger below 1 ms is raised to 1 ms. With a partition_key, every record of the handle goes to that key. Without one, each flushed batch is balanced across partitions. Rust needs the agent feature, because the linger timer runs on tokio.
| Setting | Default | Rust | Python | TypeScript |
|---|---|---|---|---|
| Record limit | 512 | max_records(n) | max_records= | maxRecords(n) |
| Byte limit | 1 MiB | max_bytes(n) | max_bytes= | maxBytes(n) |
| Linger | 5 ms | linger(Duration) | linger_ms= | linger(ms) |
| Partition key | None | partition_key(key) | partition_key= | partitionKey(key) |
send() queues one payload with headers. Rust always takes a headers map, which can be empty. Python and TypeScript make headers optional. It flushes inline when a size limit trips, so backpressure lands on the sender. flush() sends what is queued, and close() flushes and stops the timer. Flushes run one at a time, so batches reach the log in queue order. After close(), a send raises InvalidError in Python and TypeScript. Rust's close() takes the producer by value. In TypeScript, await using closes the producer at the end of its scope.
A failed timer flush never stops the timer. Its failure is kept until the next send(), flush(), or close() reports it. A send() that finds a kept failure does not queue its own record. It returns the kept failure with that record added to the unconfirmed records. flush() and close() report the kept failure after they drain the queue. When several flushes fail before you ask, you get one publish failure that lists the records of every failed batch.
await using batching = topic
.batching()
.maxRecords(512)
.maxBytes(1_048_576)
.linger(5)
.partitionKey("node-7")
.build()
const text = new TextEncoder()
await batching.send(text.encode("reading-0"))
await batching.send(text.encode("reading-1"))
await batching.flush()use laser_sdk::stream::Headers;
use std::time::Duration;
let batching = topic
.batching()?
.max_records(512)
.max_bytes(1_048_576)
.linger(Duration::from_millis(5))
.partition_key("node-7")
.build();
batching.send(b"reading-0".to_vec(), Headers::new()).await?;
batching.send(b"reading-1".to_vec(), Headers::new()).await?;
batching.flush().await?;
batching.close().await?;batching = topic.batching(
max_records=512,
max_bytes=1_048_576,
linger_ms=5,
partition_key="node-7",
)
await batching.send(b"reading-0")
await batching.send(b"reading-1")
await batching.flush()
await batching.close()When a publish fails
A publish that gives up reports what it left behind. committed lists the ranges the server confirmed, each with the stream and topic IDs, the partition, and the base offset. unconfirmed holds the records without a confirmation. Do not publish the confirmed ranges again. An unconfirmed record can already exist on the server if its acknowledgement was lost. Unconfirmed records keep the message IDs that the attempts used.
A direct producer splits a large batch into requests of batch_length records. A failure reports the confirmed requests and every record from the failed request on. Inspect the report before you retry.
| Rust | Python | TypeScript | |
|---|---|---|---|
| Failure | LaserError::PublishFailed | PublishFailedError | PublishFailedError |
| Confirmed ranges | committed | committed | committed |
| Unconfirmed records | unconfirmed, a list of IggyMessage | unconfirmed, a list of IggyMessage | unconfirmed |
| Target | stream and topic | stream and topic | stream and topic |
| Original error | publish_cause() | __cause__ | publishCause(), or the free publishCause(error) |
| Retryable | is_retryable() | retryable | isRetryable(error) |
The retryable check answers for the original error. To resend the unconfirmed records with their original message IDs, pass them back to topic.batch(..). A server that deduplicates by message ID then drops any record that already landed. The records carry no routing, so the resend does not keep the original key or partition.
import { PublishFailedError } from "@laserdata/laser-sdk"
try {
await producer.sendBatch(payloads)
} catch (error) {
if (!(error instanceof PublishFailedError)) throw error
console.error(
`${error.committed.length} range(s) confirmed, ${error.unconfirmed.length} record(s) unconfirmed`,
error.publishCause()
)
await topic.batch(error.unconfirmed)
}use laser_sdk::prelude::LaserError;
if let Err(error) = producer.send_batch(batch).await {
eprintln!("cause: {}", error.publish_cause());
if let LaserError::PublishFailed(failure) = error {
eprintln!(
"{} range(s) confirmed, {} record(s) unconfirmed",
failure.committed.len(),
failure.unconfirmed.len()
);
topic.batch(failure.unconfirmed, None).await?;
}
}import laser_sdk as ls
try:
await producer.send_batch(values)
except ls.PublishFailedError as error:
print(
f"{len(error.committed)} range(s) confirmed, "
f"{len(error.unconfirmed)} record(s) unconfirmed",
error.__cause__,
)
await topic.batch(error.unconfirmed)Consumers, offsets, and commit policies
topic.consumer(name, partition) builds a named reader for one partition. topic.consumer_group(group) returns a group handle, and group.consumer() builds a reader that joins the group. The server gives each partition to one member and rebalances as members join and leave. In Python, partition is a keyword that defaults to 0. In TypeScript, topic.consumer(name, partitionId, options) builds the named reader and await topic.consumerGroup(group).consumer(options) builds the group reader.
Every consumed message carries its log position in position, a partition and an offset. Read the offset as message.position.offset. The message also keeps the exact Iggy headers, the message ID, the checksum, and the timestamps. A malformed header block in TypeScript marks the record instead of failing the whole poll.
Delivery is at least once. A crash between delivery and commit delivers the message again, so handlers must be safe to repeat.
Two group paths
A group consumer reads one of two ways. On a server that resolves consumer group policies, such as Laser Stack and LaserData Cloud, it reads through the group and acknowledges with fenced group commits. This page calls that the group-policy path. On plain Apache Iggy it is the native Iggy group consumer, called the native path here. A group-policy consumer always joins its group, and all three SDKs refuse auto_join_group set to false on that path (TypeScript autoJoinGroup). The group-policy path is also how Filters run.
create_group and auto_join_group default to true (Python create_group= and auto_join_group=, TypeScript createGroup and autoJoinGroup). The server stores offsets by consumer name or group. Reconnect with the same name to resume from the stored offset.
Where a consumer starts
start_at(..) picks the first position when no stored offset exists. The default is Next.
- Rust takes
ConsumerStart::First,Last,Next,Offset(n), orTimestampMicros(t). - TypeScript takes
startAtwith{ kind: "first" | "last" | "next" }, or{ kind: "offset" | "timestampMicros", value }wherevalueis abigint. - Python takes
polling="first" | "last" | "next", oroffset=ortimestamp_micros=.
allow_replay() in Rust, allow_replay=True in Python, and allowReplay: true in TypeScript permit reads at or below a stored offset. Use it to replay a group from First. Without it, a native consumer skips records it already consumed. A group-policy consumer honors its start position as given.
Next does not skip the first record of a fresh consumer. With no stored offset, it starts at the first retained record. With a stored offset of 0, it starts at offset 1. Storing offset 0 acknowledges that record. It is not an empty starting state.
Commit and offset control
commit(message) acknowledges a processed record. Native consumers also expose store_offset(offset, partition) and delete_offset(partition) for direct offset control (TypeScript storeOffset(offset, partitionId?) and deleteOffset(partitionId?), Python store_offset(offset, partition=) and delete_offset(partition=)). Group-policy consumers refuse those two calls, because their acknowledgments must pass the group and policy checks.
last_consumed_offset(p) and last_stored_offset(p) read local progress (TypeScript lastConsumedOffset(p) and lastStoredOffset(p), awaited in Python). They are diagnostics, not a resume point. On the native path in Rust and Python, the Iggy cache can hold an initial zero before any commit. On the group-policy path, last_stored_offset reports the last acknowledged offset. Resume with the Next start instead of computing a position from them. stored_offset(p) (TypeScript storedOffset(p)) reads the stored offset and the partition's current offset from the server, or returns nothing when the server stores none yet. TypeScript refuses it on an unnamed partition consumer.
A purge restarts every partition at offset 0. An open native consumer can keep its old position and skip the new records, so shut it down and build a new one with the same name or group. Group-policy consumers detect the changed history, reset partition progress, and report source_changed before the next read resumes. A group handle built from a numeric ID refuses a recreated stream or topic and needs a new reader.
Commit policies
The commit policy decides when the consumer stores offsets on the server:
| Policy | Stores the offset |
|---|---|
Polling (default) | Native consumers: the end of the polled batch, 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) / IntervalOrEvery(d, n) | The interval or the event, whichever comes first |
Disabled | Never. You call commit(message) yourself |
On the native path the default policy can skip unprocessed records after a crash, because it commits before delivery. Group-policy consumers acknowledge the delivered prefix when consumption continues or shuts down, so a crash delivers the current batch again. Use Disabled and commit after successful processing when only finished work may advance. Native consumers send an offset store, and group-policy consumers send a fenced acknowledgment.
await using consumer = await topic.consumerGroup("workers").consumer({
batchLength: 100,
commitPolicy: { kind: "disabled" },
startAt: { kind: "first" },
allowReplay: true,
pollIntervalMs: 5
})
for await (const message of consumer) {
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()To commit automatically instead, drop the manual commits and set a policy: commit_policy(CommitPolicy::IntervalOrEach(Duration::from_secs(1))) in Rust, auto_commit="each", commit_interval_ms=1000 in Python, or commitPolicy: { kind: "intervalOrEach", intervalMs: 1_000 } in TypeScript. In TypeScript, consumer.stream({ signal }) takes an AbortSignal. Aborting it stops the loop with a CancelledError.
Waiting with a timeout
next_within(timeout) returns the next record, or fails when the time runs out. Rust takes a Duration, and Python and TypeScript take milliseconds. A timeout is a typed error: LaserError::Timeout in Rust and TimeoutError in Python and TypeScript. After shutdown, the call fails with an invalid-state error.
TypeScript keeps a poll that outlives the timeout. The next call resumes that poll and can receive its late record. A timeout never starts overlapping polls and never acknowledges a record that the caller did not receive.
import { TimeoutError } from "@laserdata/laser-sdk"
try {
const message = await consumer.nextWithin(5_000)
handle(message)
} catch (error) {
if (!(error instanceof TimeoutError)) throw error
// nothing arrived in five seconds
}match consumer.next_within(Duration::from_secs(5)).await {
Ok(message) => handle(&message)?,
Err(LaserError::Timeout(_)) => {
// nothing arrived in five seconds
}
Err(error) => return Err(error),
}import laser_sdk as ls
try:
message = await consumer.next_within(5_000)
handle(message)
except ls.TimeoutError:
pass # nothing arrived in five secondsDifferences by language
All three SDKs share these consumer defaults: batch_length 1000, the Polling policy, the Next start, a 1 second wait before a failed poll is retried, and group creation and join on build. All three also set the failed-poll wait and the init retries: Rust polling_retry_interval and init_retries, Python polling_retry_interval_ms, init_retries, and init_retry_interval_ms, TypeScript pollingRetryIntervalMs and initRetries: { retries, intervalMs }. These parts differ:
- Rust sets one of the ten
CommitPolicyvariants withcommit_policy(..). Withoutpoll_interval, a native consumer polls again at once after an empty poll. A group-policy consumer waits 250 ms. - Python picks a policy with
auto_commit:disabled,interval,polling(the default),all,each, orevery.commit_interval_msdefaults to 0, which keeps the plain policy. A positive value selects the interval variant, such asIntervalOrEach.intervalneeds a positivecommit_interval_ms, andeveryneedscommit_every.poll_interval_msfollows the Rust default. The consumer initializes on first read. Callawait consumer.init()when startup must fail before the service accepts work. - TypeScript sets
commitPolicywith the same ten variants as{ kind, intervalMs, count }objects, such as{ kind: "disabled" }or{ kind: "every", count: 10 }.pollIntervalMsdefaults to 0 on native consumers, which poll again without a delay.
shutdown() stops polling and leaves the group. Automatic policies store the handled prefix first. Native polling commits before delivery, so its shutdown does not prove that every record was processed. If a native consumer delivered offset 0 on a partition, automatic shutdown stores offset 0 explicitly unless an earlier store already covers it. With Disabled, progress follows your commits.
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
Filters let one shared topic deliver a different subset to each consumer group. The server selects records before they cross the network. Configure the policy through topic.consumer_group(name).filter(). Normal group consumers and group.reader() apply it, and the group's progress still covers every record the server scanned. Raw Iggy consumers ignore the filter and receive every record, so do not mix them with Laser consumers on one group.
Apache Iggy access
Rust's topic.iggy_producer(), topic.iggy_consumer(..), and topic.iggy_consumer_group(group) return native Iggy builders on the same connection, and laser.client() returns the client. Use them for settings that Laser SDK does not expose. TypeScript exposes the underlying client as laser.client. Python has no raw Iggy accessor.
Durability
Publish completion follows the topic's durability policy. A local queue that accepted a record does not prove the server stored it, and neither does a background send that returned. Message durability and consumer offset durability are independent. See Topic durability and Acknowledgements and offset commits.
Producer statistics
Producers in all three SDKs can report statistics over a separate observer connection. Set the variables before you create producers. A telemetry error never fails a publish.
| Setting | 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, no headers or framing |
| Confirmed records and bytes | Fully successful synchronous send calls. Unavailable for background sends |
| Failed calls | Send calls that returned an error |
| Last success | The last successful application call |
| Retries | Unavailable when the native client does not expose its attempts |
| Latency p50, p99, p99.9 | Publish-call time including retry delays, as power-of-two bucket upper bounds. p99 needs 100 calls, p99.9 needs 1000 |
- A failed batch can have committed a prefix that the confirmed counters leave out. The counters are SDK reports. Do not use them as delivery totals, exactly-once evidence, usage metering, or authorization input.
- A value outside its range is clamped to the range. A report expires after three intervals. Handles past the limit publish normally and are left out. A producer ID stays stable for the life of its handle, across reconnects.
- A Laser built from an injected client that cannot open a second authenticated connection does not report. Reports also need a server that answers the AGDX hello, so plain Apache Iggy gets none.
- 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.
Key operations
| Operation | What it does |
|---|---|
stream(name).topic(name) | Address a topic on any stream |
topic.ensure(partitions) | Create the topic if it does not exist, with a fixed partition count and no message expiry |
publish().payload(bytes) | Append one raw message |
.json(v) / .msgpack(v) | Encode the body before publish |
.avro(schema, schema_id, v) | Encode a body as an Avro datum under a registered writer schema. Rust needs the schema-codecs feature |
.claim_check(store, threshold) | Store a large body in a blob store and publish a reference (TypeScript claimCheck). Rust needs the agent feature |
.partition_key(key) | Route by key for per-key order |
.index(key, value) | Stamp an indexed text header, so a view can query it |
.header(key, value) | Attach a text header that stays on the record and is not queryable. Use topic.producer() for typed headers |
publish_batch() | Collect several messages into one request |
topic.json(..) / topic.cbor(..) | A typed handle. publish encodes, publish_batch encodes many, records(name) opens a typed cursor |
topic.schema(id) | A typed handle bound to a registered schema. Async in all three SDKs. Needs Laser Stack or LaserData Cloud, and Rust needs the schema-codecs feature |
topic.producer() | A long-lived producer with batching, linger, and retries. See Producing at volume |
topic.batching() | A client-side accumulator with flush() and close(). See Batching producers |
topic.consumer(..) / topic.consumer_group(name).consumer() | A live consumer, partitioned across group members |
consumer.next() / next_within(timeout) | Wait for the next record, optionally with a typed timeout |
consumer.commit(message) | Commit after the record is handled |
store_offset(..) / delete_offset(..) | Native consumers only |
topic.replay() | A raw cursor from offset 0, or from saved offsets, across every partition |
An indexed header wins over a value that the projection schema extracts for the same field. A record with no indexed fields from headers or the schema produces no view row. Queries and views explains extraction and projection registration.
A typed publish through topic.schema(id) encodes Avro and JSON Schema bodies. For Protobuf, Rust and Python publish the encoded bytes with the schema ID, and TypeScript encodes the body.
BatchPublishRequest.add_record in Python takes per-record metadata. Batch defaults fill what the record leaves unset, and batch index entries and headers merge with the record's own, which win. Its logical_schema_fingerprint= parameter takes 32 bytes and needs Arrow content. TypeScript exposes the same metadata through Record.logicalSchemaFingerprint(bytes). Use the Arrow IPC helpers when the SDK must check metadata and payload length.
Sessions in depth
Every session option, default, limit, and read surface: lifecycle, leases, layouts, agents, submission, operator control, budgets, pause, state, replay, and managed reads
Filters in depth
Edge cases, batch semantics, progress, revisions, errors, permissions, limits, and payload formats for consumer group filters