LaserData Cloud
Laser SDK

Consumer Filters

Filter records before they cross the network. Receive only what each consumer needs.

Receive only the records you need. A consumer filter selects records on the streaming server, before they cross the network. Change data capture (CDC) records describe changes published by a producer. The reader asks for a page of one partition. The server scans it and returns only the records the filter selects. Each record keeps its original offset, id, timestamp, headers, and payload bytes. Everything else stays on the server.

98.5% less payload transfer in the complete CDC example. The reader receives 4 of 240 records and 424 of 27,953 payload bytes. Rust, Python, and TypeScript produce the same result.

At scale the picture holds. The Frostline example publishes ten million change records, 6.30 GB of payload, and reads them with four filtered groups. The groups receive 777.64 MB where reading the full feed four times would move 25.21 GB, 96.9% of payload avoided. Counting both TCP directions with strace gives 1.11 GB against 31.40 GB, 96.5% less TCP application traffic. These counts exclude TCP/IP headers and retransmissions. The same repository holds the benchmark, p50/p99/p99.9 tables, server CPU, and memory results. Traffic reduction is measured separately from latency and CPU cost. Header predicates can avoid payload decoding. Payload predicates add server work while reducing downstream bytes and application work.

Built for

Use consumer filters for CDC consumers that care about one kind of change, for alert routing, and for backfills that must skip most of a topic.

How it works

A change feed carries every change to every table. Without a filter, a consumer downloads and decodes all of it just to drop most records. With a filter, the server does that work next to the data, and the log stays the single source of truth. Savings depend on the matching fraction of payload bytes in your feed.

One topic can serve many selective consumers without sending every record to each one. Consumers can follow different subsets of the same topic and its partitions. You don't need an intermediate topic just to deliver each subset.

For example, suppose ten consumers each need 5% of a topic's payload bytes:

Payload from one 1 GiB source rangeRead everything, filter in each applicationFilter on the server
One consumer1 GiB51.2 MiB
Ten consumers combined10 GiB512 MiB
Payload transfer saved0%95%

This is an illustration of payload transfer, excluding protocol overhead. The broker still examines records for each group. Matching records retain their full payload. Applications avoid receiving and decoding rejected payloads.

A filter belongs to a consumer group. The hierarchy is Laser → stream → topic → consumer group → filter. You create the group once, give it a filter, and every Laser SDK consumer of that group runs it. A consumer names only the group. A group with no filter receives every record. Changing a group's filter is a new group, so shared offsets never silently cover a different selection.

Raw Iggy polling ignores the filter. topic.iggy_consumer_group(..), the Iggy CLI, and any native Iggy client poll the same group with standard semantics and receive every record. The filter is a delivery policy for Laser SDK consumers, not an access-control boundary, and a native source-read grant still permits native consumption. Don't mix the two on one group: a native poll stores the group's offset past records the filtered consumers never saw, and they resume after it.

A filter is declarative. It never runs your code. What it can test:

TestOperators
Payload field comparisoneq, ne, lt, lte, gt, gte, in, contains, prefix
Text matching on a field or a headerequals, prefix, suffix, contains, glob, regex, each optionally case-insensitive
Typed header comparisonthe same operators, by exact header key, without decoding the payload
Field presencepresent, absent
Compositionall, any, not
Explicit coercioncompare a timestamp string as an instant or a decimal string as a number

The payload codec is part of the filter. JSON, CBOR, Avro, and Protobuf support field filtering. headers_only never decodes the payload, so it works with any format. Avro and Protobuf require registered writer schemas.

Truth has three values. Most comparisons with a missing field or a type mismatch return unknown, and not keeps unknown. An unknown result selects nothing, unless the filter's mismatch policy asks to see type mismatches (see Edge cases). Null comparisons are the exception: eq null matches a missing or explicit-null field, and ne null excludes both. Use present and absent when the distinction matters.

A payload that cannot be decoded is a fault, which is separate from truth. The filter's fault policy decides what happens:

PolicyWhat happens to the record
stop (default)The page stops before it. The reader raises the fault and reads that partition again after its idle interval, so the block clears once the record ages out of the partition. To move past it sooner, create a reader with an explicit start position or another policy
passDelivered and marked as not evaluated
dropSkipped and included in the examined count

Edge cases: foreign and mismatched records

Two record policies decide what happens to records a filter cannot judge on its own terms. Both default to reject.

PolicyCoversreject (default)pass
foreign_policyA record whose agdx.ct header names another codec than the filter's, or whose writer schema is missing, not listed by the filter, or of another schema familySkipped. Never decoded and never a faultDelivered and marked as not evaluated
mismatch_policyA record where a predicate's value exists but has a type the predicate cannot compare, for example xyz == "Abc" against "xyz": 1, an object, or an array, and the filter could not decide otherwiseSkippedDelivered and marked as not evaluated

A missing or null field is not a mismatch. It is plain unknown and never delivered by the mismatch policy. A record without an agdx.ct header is decoded with the filter's codec. Content type any and codes a server does not know also decode as usual. A payload that is broken in the filter's own codec is still a fault and follows the fault policy.

Records delivered under pass are listed in the page's unevaluated offsets, so a consumer can quarantine them, log them, or apply its own rules.

const strict = ConsumerFilter.json(FilterExpr.pred("xyz", "eq", "Abc"))
const seeEdgeCases = ConsumerFilter.withMismatchPolicy(strict, "pass")
let see_edge_cases = ConsumerFilter::json(FilterExpr::pred("xyz", CmpOp::Eq, "Abc"))
    .with_mismatch_policy(RecordPolicy::Pass);
strict = ls.ConsumerFilter.json(ls.FilterExpr.pred("xyz", "eq", "Abc"))
see_edge_cases = strict.with_mismatch_policy("pass")

Every filter has a digest, a SHA-256 hash of its canonical form. The digest names exactly what the filter selects, so two filters with the same digest behave the same way. A group remembers the digest it was configured with.

Define a filter

Build the expression, then wrap it in the codec it reads. This filter follows a satellite fleet. It selects a producer-reported mode change to safe, a decommission, or a telemetry report of safe mode. The complete examples publish typed records, a serde model in Rust, dataclasses in Python, and a discriminated union in TypeScript, and decode what arrives back into those types.

import { ConsumerFilter, FilterExpr } from "@laserdata/laser-sdk"

const satellites = () => FilterExpr.pred("table", "eq", "satellites")
const safeMode = ConsumerFilter.json(
  FilterExpr.any([
    FilterExpr.all([
      satellites(),
      FilterExpr.pred("op", "eq", "u"),
      FilterExpr.pred("changed", "contains", "mode"),
      FilterExpr.pred("after.mode", "eq", "safe")
    ]),
    FilterExpr.all([satellites(), FilterExpr.pred("op", "eq", "d")]),
    FilterExpr.all([
      FilterExpr.pred("event", "eq", "satellite.telemetry_changed"),
      FilterExpr.pred("fields.mode", "eq", "safe")
    ])
  ])
)
use laser_sdk::filters::{ConsumerFilter, FilterExpr};
use laser_sdk::query::CmpOp;

let satellites = || FilterExpr::pred("table", CmpOp::Eq, "satellites");
let safe_mode = ConsumerFilter::json(FilterExpr::any([
    FilterExpr::all([
        satellites(),
        FilterExpr::pred("op", CmpOp::Eq, "u"),
        FilterExpr::pred("changed", CmpOp::Contains, "mode"),
        FilterExpr::pred("after.mode", CmpOp::Eq, "safe"),
    ]),
    FilterExpr::all([satellites(), FilterExpr::pred("op", CmpOp::Eq, "d")]),
    FilterExpr::all([
        FilterExpr::pred("event", CmpOp::Eq, "satellite.telemetry_changed"),
        FilterExpr::pred("fields.mode", CmpOp::Eq, "safe"),
    ]),
]));
import laser_sdk as ls

def satellites() -> ls.FilterExpr:
    return ls.FilterExpr.pred("table", "eq", "satellites")

safe_mode = ls.ConsumerFilter.json(
    ls.FilterExpr.any([
        ls.FilterExpr.all([
            satellites(),
            ls.FilterExpr.pred("op", "eq", "u"),
            ls.FilterExpr.pred("changed", "contains", "mode"),
            ls.FilterExpr.pred("after.mode", "eq", "safe"),
        ]),
        ls.FilterExpr.all([satellites(), ls.FilterExpr.pred("op", "eq", "d")]),
        ls.FilterExpr.all([
            ls.FilterExpr.pred("event", "eq", "satellite.telemetry_changed"),
            ls.FilterExpr.pred("fields.mode", "eq", "safe"),
        ]),
    ])
)

Group reads are part of the default streaming SDK. The Rust filters Cargo feature, also included by managed, adds only the local evaluator and the reader's optional local guard:

laser-sdk = { version = "0.5", features = ["filters"] }

Create Laser from a connection string, so the SDK can open authenticated connections to the serving nodes. A caller-supplied raw client without connection settings supports local diagnostics only, without acknowledgments.

Configure once, consume by group ID

Create the group with its filter once. Consumers then name only the group. The server runs the group's filter, or none, on every read. Several instances of the same group share partition assignments and progress. Different groups keep independent progress and can each receive their own copy of the same matching records.

const topic = laser.stream("orbit").topic("fleet_changes")
const desk = topic.consumerGroup("anomaly-desk")

// One-time setup. Omit filter to create a group that reads everything.
const info = await desk.create({ filter: safeMode })
const groupId = info.id

// Consumer instances only need the topic and group ID.
const consumer = await topic.consumerGroupId(groupId).consumer({
  autoCommit: false,
  batchLength: 100
})
try {
  for await (const record of consumer) {
    // FleetChange is the record union of the complete example.
    const change = JSON.parse(new TextDecoder().decode(record.payload)) as FleetChange
    console.log(record.partitionId, record.offset, change)
    await consumer.commit(record)
  }
} finally {
  await consumer.shutdown()
}
use laser_sdk::prelude::{CommitPolicy, LaserError};

let topic = laser.stream("orbit").topic("fleet_changes");
let desk = topic.consumer_group("anomaly-desk");

// One-time setup. Omit .filter(..) to create a group that reads everything.
let info = desk.create().filter(safe_mode.clone()).build().await?;
let group_id = info.id;

// Consumer instances only need the topic and group ID.
let mut consumer = topic
    .consumer_group_id(u64::from(group_id))
    .consumer()
    .batch_length(100)
    .commit_policy(CommitPolicy::Disabled)
    .build()
    .await?;
while let Some(record) = consumer.next().await {
    let record = record?;
    // FleetChange is the serde model of the complete example.
    let change: FleetChange = record.json()?;
    println!("{} {} {change:?}", record.partition_id, record.position.offset);
    consumer.commit(&record).await?;
}
consumer.shutdown().await?;
topic = laser.stream("orbit").topic("fleet_changes")
desk = topic.consumer_group("anomaly-desk")

# One-time setup. Omit filter to create a group that reads everything.
info = await desk.create(filter=safe_mode)
group_id = info["id"]

# Consumer instances only need the topic and group ID.
consumer = topic.consumer_group_id(group_id).consumer(
    auto_commit="disabled", batch_length=100
)
try:
    async for record in consumer:
        print(record.partition_id, record.offset, record.json())
        await consumer.commit(record)
finally:
    await consumer.shutdown()

create is idempotent. An existing group is kept, and the same definition keeps its binding. A different definition on a group that already runs one returns conflict. Create another group instead, so two selections never share one set of offsets. You can also create a plain group first and call group.filter().configure(filter) later. The definition is saved as the group's own filter in the catalog, with a filter ID and revision 1, and the group is bound to it in one step.

Native group creation and catalog configuration are two operations. If configuration fails after the group exists, Rust returns LaserError::ConsumerGroupSetup, TypeScript throws ConsumerGroupSetupError, and Python raises FilterError, each with the group's ID and identity. The group is left as it is, unbound. Pass an operation_id you recorded first to resume the same configuration after a crash and read its first outcome.

Complete examples: Rust, Python, and TypeScript.

What a batch of 100 means

A normal consumer with batch length 100 examines at most 100 source records per partition request and delivers zero through 100 records. If seven match, you receive seven. If none match, the SDK still keeps the scan position and continues through the backlog. An empty poll does not mean the topic is empty.

batch_length is both the page size and the scan budget on the normal consumer, so one poll never costs the server more than one batch of source records. A selective filter over a long backlog therefore takes many polls to reach the first match. When you want the server to scan further per request, use the group reader.

The group reader separates the two limits. count bounds delivered matches and max_examined bounds scanned source records, within the server's own limits. It returns pages, acknowledges explicitly, and reads every record when the group is unbound. Use next_record for a record loop or next_page for batch handling. TypeScript uses nextRecord and nextPage.

const reader = await desk
  .reader()
  .start({ kind: "first" })
  .count(100)
  .maxExamined(1000)
  .build()
try {
  const page = await reader.nextPage({ timeoutMs: 15_000 })
  for (const record of page.records) {
    console.log(record.partitionId, record.offset, record.json())
  }
  await reader.ackPage(page)
} finally {
  await reader.close()
}
use laser_sdk::filters::FilteredStart;

let mut reader = desk
    .reader()?
    .start(FilteredStart::First)
    .count(100)
    .max_examined(1000)
    .build()
    .await?;
let page = reader.next_page().await?;
for record in &page.records {
    let change: FleetChange = record.json()?;
    println!("{} {} {change:?}", record.partition_id, record.offset);
}
reader.ack_page(&page).await?;
reader.close().await?;
reader = await desk.reader(start="first", count=100, max_examined=1000)
try:
    page = await reader.next_page()
    for record in page.records:
        print(record.partition_id, record.offset, record.message.json())
    await reader.ack_page(page)
finally:
    await reader.close()

A page is one bounded filtered poll result, with matching records and scan progress. It is not a separate transport or a CDC-only storage format. The AGDX command uses standard Iggy transport and leaves stored records unchanged. For batch handling, process every returned record before ack_page (ackPage in TypeScript). For record handling, call ack(record) after processing. The SDK advances stored progress only through the completed prefix.

Raw Iggy polling keeps its standard behavior and does not apply the group's filter, see How it works.

Safe progress and policy changes

A fresh group using Next starts at the first retained record, including offset 0. After offset 0 is acknowledged, Next resumes at offset 1. An unset checkpoint and a stored checkpoint of zero are different states. Don't initialize a new group by storing zero.

Acknowledgments cover completed work, including records the server skipped. A primary page that examines records reports a safe offset, the last record it finished classifying. When you acknowledge a page, the server stores that offset, including the records it skipped. A consumer that matches one record in a million still moves forward. A page that examines offsets 0 through 99 and returns matches at 10 and 20 can store progress through 99 once both are done. It cannot cross an earlier unhandled delivery.

On the normal consumer the commit policy decides when the handled prefix is stored. Polling (the default) stores it at the next poll and at shutdown, Interval on a timer, Each, Every(n) and All after deliveries, and Disabled only when you call commit(record). TypeScript maps autoCommit: true to the polling behavior. Nothing is stored before a record reaches your code, so a crash redelivers the batch you were handling instead of skipping it. Use Disabled plus commit when processing must finish before progress is stored.

The server fences every acknowledgment by source incarnation, group identity, execution mode and policy generation. An old unfiltered page cannot commit after a filter is configured. An old filtered page cannot commit after the group is released. A purge or a recreated topic invalidates old delivery handles. A revoked partition can still drain accepted pages while its commit fence permits it.

A permission change takes effect inside a page. The server checks source read permission before every read round. A revocation between rounds fails the request and returns no partial page. The group's policy resolves once when the request starts, so a catalog change applies from the next request. An acknowledgment checks offset store permission when it runs.

The SDK rejects a page or record from another reader and validates response identity and progress before exposing it. It reads ahead and lets you finish records in any order. It stores only the offset up to which every returned record is done. Pages with no matches store their scanned range on their own.

A reader keeps at most 1024 unacknowledged record-bearing pages per partition and stops reading that partition until you acknowledge, so read-ahead memory stays bounded. Change the bound with max_unacked_pages in Rust and Python or maxUnackedPages in TypeScript. Empty pages don't count.

Reads go to the partition primary by default. That is the only mode that can acknowledge. The local read mode serves diagnostics from whatever replica the connection reaches. A local page offers no acknowledgment offset, because a lagging replica can still hold records from before a purge.

A group reader joins over its own connection, so two readers of one Laser are two members, and a reader dropped without close leaves the group with its connection. A partition the group hands a reader later resumes after the group's stored offset, whatever start the reader was built with, so the previous owner's unacknowledged backlog is read, not skipped. A partition that fails waits one idle interval while the other partitions keep reading.

Pausing a revision stops new reads. Work already delivered can still drain under the acknowledgment rules. A failed catalog lookup, an unavailable plane, or a paused filter never silently becomes an unfiltered read.

Releasing or deleting a filter does not stop its consumers. From their next poll they receive every record. Records delivered under the old filter but not yet stored are delivered again from the stored offset, nothing is skipped. Only an explicit commit or ack of a record read under the old filter fails, with conflict, because that offset was not stored.

Checkpoints inside a page

Persist the application checkpoint before acknowledging its offset. Rust and Python expose ack_through(record). TypeScript exposes ackThrough(record). The call marks all preceding returned records on that partition as handled and stores the completed prefix. Process that full prefix first. Later records in the same page remain pending. Use this when a window receipt or transaction ends inside a page. Use ack_page or ackPage when the whole page is complete.

Cluster freshness

Configuration outcomes carry a durable operation receipt and a control-log position. The group handle remembers its last configuration, including across a cloned handle, and every read built from it carries that position. The serving plane confirms that operation before it answers, so an application never sees its own configuration as absent through another node.

A read from a group nobody configured recently can hit a cached unbound answer. IGGY_PLANE_FILTERS_UNBOUND_CONFIRM_MS bounds how long a broker shard reuses that answer before it asks the plane for a lookup confirmed against the control-log head. The default is 1000 ms. Zero confirms every read. A filter configured through another node takes effect within that window at the latest. Don't treat it as an instantaneous cross-node guarantee.

Acknowledgments confirm the authoritative policy again. A node that cannot prove its catalog is current returns the retryable catalog_unavailable. Cluster membership, primary routing and replication keep their native Iggy behavior.

What a change record proves

Filters evaluate each record independently. They don't compare it with an earlier version of the entity. In the example, the producer supplies changed, a list of columns that changed. The snapshot branch requires both changed containing mode and after.mode equal to safe.

A battery update from a satellite already in safe mode does not match that branch. A values-only filter selects it. Choose that filter when you need current state rather than evidence of a transition.

The telemetry branch selects a report of safe mode. Repeated or forced telemetry reports can match even when the mode did not change. The decommission branch selects deletions independently of mode. Keep these producer semantics explicit when adapting the example to your feed.

Test and preview

Check the group's filter before you depend on it. Neither call joins the group or stores an offset.

  • test judges one sample record you supply and explains every predicate.
  • preview judges a range of stored records in one partition and lists each verdict.
const tested = await desk.filter().test(JSON.stringify(sample))
console.log(tested.explanation.verdict)

const preview = await desk.filter().preview(0, { maxRecords: 10 })
console.log(preview.examined, preview.matched, preview.stop)
let tested = desk.filter().test(sample, Vec::new()).await?;
println!("{}", tested.explanation.verdict);

let preview = desk.filter().preview(0).await?.max_records(10).send().await?;
println!("{} {} {}", preview.examined, preview.matched, preview.stop);
tested = await desk.filter().test(json.dumps(sample))
print(tested["explanation"]["verdict"])

preview = await desk.filter().preview(0, max_records=10)
print(preview["examined"], preview["matched"], preview["stop"])

test also takes typed headers, so a headers-only filter can be checked the same way. The local evaluator runs the same verdict without a server: ConsumerFilter.evaluate(payload) in Python, CompiledFilter.compile(filter).evaluate(record) in TypeScript, and CompiledFilter from laser_sdk::filters in Rust.

Filter on headers only

Filter opaque payloads without decoding them. A headers_only filter reads typed user headers and never touches the payload. Use it when producers stamp a routing header and the payload is Protobuf, Avro, or any other format. The pager group below receives only critical alert frames.

import { HeaderValue } from "@laserdata/laser-sdk"

const alerts = laser.stream("orbit").topic("fleet_alerts")
await alerts.producer().send(frame, { headers: { priority: HeaderValue.uint8(2) } })

const pager = alerts.consumerGroup("pager")
await pager.create({
  filter: ConsumerFilter.headersOnly(FilterExpr.header("priority", "eq", 2))
})
const reader = await pager.reader().start({ kind: "first" }).build()
const record = await reader.nextRecord({ timeoutMs: 15_000 })
console.log(record.offset, record.payload.byteLength)
await reader.ack(record)
await reader.close()
use laser_sdk::iggy::prelude::{HeaderKey, HeaderValue};
use laser_sdk::stream::ProducerMessage;

let alerts = laser.stream("orbit").topic("fleet_alerts");
let producer = alerts.producer().build().await?;
let message = ProducerMessage::new(frame)
    .header(HeaderKey::try_from("priority")?, HeaderValue::from(2_u8));
producer.send_message(message).await?;

let pager = alerts.consumer_group("pager");
pager
    .create()
    .filter(ConsumerFilter::headers_only(FilterExpr::header("priority", CmpOp::Eq, 2_i32)))
    .build()
    .await?;
let mut reader = pager.reader()?.start(FilteredStart::First).build().await?;
let record = reader.next_record().await?;
println!("{} {}", record.offset, record.message.payload.len());
reader.ack(&record).await?;
reader.close().await?;
alerts = laser.stream("orbit").topic("fleet_alerts")
producer = alerts.producer(partitions=1)
await producer.send(frame, headers={"priority": ("uint8", 2)})

pager = alerts.consumer_group("pager")
await pager.create(
    filter=ls.ConsumerFilter.headers_only(ls.FilterExpr.header("priority", "eq", 2))
)
reader = await pager.reader(start="first")
record = await reader.next_record()
print(record.offset, len(record.message.payload))
await reader.ack(record)
await reader.close()

Header keys are custom text names matched exactly, including dots. Values can be booleans, signed or unsigned integers through 64 bits, floating-point numbers, or strings. Choose an exact-width Iggy header when space matters: uint8 uses one value byte. Numeric 2 and text "2" are different values. Use the streaming producer for typed headers. The publish().header(..) convenience accepts strings.

frame above is your encoded payload. Header-only selection works with JSON, Avro, Protobuf, or other bytes, because it does not inspect the payload. Non-text keys and 128-bit numeric comparisons are outside the current filter contract.

One log, many kinds of events

Put every producer's events in one ordered log and let each group select its subdomain on the server. One partition keeps a single global order. More partitions add parallel group workers. Events can differ in type, payload shape, envelope, and codec.

The pattern that holds up:

  1. Stamp a routing header on every record, such as event.type = "billing.invoice.v1.created", and the content type through the codec helpers, which set agdx.ct.
  2. Select the subdomain with a header text match, such as prefix billing., contains .v1., or a glob or regex. A headers_only filter never decodes a payload, so every format mixes safely, and it is the cheapest filter to run.
  3. Add payload conditions only where needed, inside all([header match, payload predicates]). Header children run first and all stops at the first no-match, so records of other kinds are rejected before any decode.
  4. Keep foreign_policy at reject. A record in another codec is then skipped instead of stalling the partition, even when its routing header is missing.
const billingV1 = ConsumerFilter.headersOnly(
  FilterExpr.headerText("event.type", "glob", "billing.*.v1.*")
)
const largeInvoices = ConsumerFilter.json(
  FilterExpr.all([
    FilterExpr.headerText("event.type", "prefix", "billing.invoice."),
    FilterExpr.pred("amount_cents", "gte", 100_000)
  ])
)
let billing_v1 = ConsumerFilter::headers_only(FilterExpr::header_text(
    "event.type",
    TextMatch::Glob,
    "billing.*.v1.*",
));
let large_invoices = ConsumerFilter::json(FilterExpr::all([
    FilterExpr::header_text("event.type", TextMatch::Prefix, "billing.invoice."),
    FilterExpr::pred("amount_cents", CmpOp::Gte, 100_000_i64),
]));
billing_v1 = ls.ConsumerFilter.headers_only(
    ls.FilterExpr.header_text("event.type", "glob", "billing.*.v1.*")
)
large_invoices = ls.ConsumerFilter.json(ls.FilterExpr.all([
    ls.FilterExpr.header_text("event.type", "prefix", "billing.invoice."),
    ls.FilterExpr.pred("amount_cents", "gte", 100_000),
]))

Text matching

KindMatches whenRelative cost
equalsthe whole value equals the patternlowest
prefixthe value starts with the patternlowest
suffixthe value ends with the patternlowest
containsthe pattern occurs anywherelow, one scan of the value
globthe whole value matches, * is any run, ? is one character, \ escapeslow to moderate with many stars
regexthe value contains a matchmoderate, linear in the value

Add case-insensitive matching to any kind. It lowercases the value first, which allocates a copy. The value must be text. A number, object, or array is a type mismatch and follows the mismatch policy.

On a JSON change record of about 600 bytes, one evaluation took about 1.3 to 1.5 microseconds for every text kind, the same as a plain equality test, because parsing the payload dominates. A header-only test took about 8 nanoseconds. Measure your own records with the SDK's filter_evaluation benchmark.

Regular expressions run in linear time on the server. A pattern cannot backtrack exponentially. Lookaround, backreferences, and inline flags such as (?i) are refused. The server's Rust engine decides pattern syntax. TypeScript performs structural validation and sends patterns to that engine. Use the case-insensitive option instead of (?i). Character classes such as \d and \w are Unicode-aware on the server. The TypeScript local guard rejects filters that contain a regex, because JavaScript's classes differ. Disable localGuard for these filters, or use Rust or Python for local verification. The Rust and Python guards run the server's evaluator. Case-insensitive regex uses the engine's Unicode folding. Other text predicates lowercase the pattern and value.

Revisions, pause and resume, A/B

A group's filter has immutable revisions. revise drafts a new revision of the group's own definition. Readers keep running the active one. A draft never replaces what a bound group runs. To run the stricter variant, create a second group with that definition. Each group has its own offsets, so one variant never skips the other's records.

set_revision_enabled(revision, false) pauses the active revision. New reads then fail with revision_disabled. Records already delivered can still be acknowledged, and test and preview still work. Set it to true to resume. The flag changes neither the definition nor its digest.

release() removes the group's policy. Its consumers then receive every record, and its filter stays saved in the catalog. delete() removes the group's own definition and every revision, and releases its binding in the same catalog operation. The group retains its previous digest restriction after either operation. Configure it again with the same digest, or create another group for a different policy. A named filter dropped from the console or over HTTP also releases its groups and leaves a tombstone that reserves the name. Past the tombstone limit the oldest tombstones are deleted. An id is never reused.

const binding = await desk.filter().get()
if (binding === undefined) throw new Error("the desk is unbound")

// Draft on the desk's own filter. The desk keeps running revision 1.
const draft = await desk.filter().revise(binding.revision, ConsumerFilter.json(transitionsOnly))
const revisions = await desk.filter().revisions({ page: 0, pageSize: 10 })

// Run the variant in its own group, with its own offsets.
const variant = topic.consumerGroup("anomaly-desk-transitions")
const created = await variant.create({ filter: ConsumerFilter.json(transitionsOnly) })

// Pause and resume the variant.
if (created.filter === undefined) throw new Error("the variant was created unbound")
await variant.filter().setRevisionEnabled(created.filter.revision, false)
await variant.filter().setRevisionEnabled(created.filter.revision, true)

// Another definition on a running group is refused.
try {
  await desk.filter().configure(ConsumerFilter.json(transitionsOnly))
} catch (error) {
  if (!(error instanceof FilterExecutionError) || error.reason !== "conflict") throw error
}

// Release: the desk receives every record again.
const released = await desk.filter().release()
let binding = desk.filter().get().await?.expect("the desk is bound");

// Draft on the desk's own filter. The desk keeps running revision 1.
let draft = desk
    .filter()
    .revise(binding.revision, ConsumerFilter::json(transitions_only.clone()))
    .await?;
let revisions = desk.filter().revisions(0, 10).await?;

// Run the variant in its own group, with its own offsets.
let variant = topic.consumer_group("anomaly-desk-transitions");
let created = variant
    .create()
    .filter(ConsumerFilter::json(transitions_only.clone()))
    .build()
    .await?;
let variant_revision = created.filter.expect("bound").revision;

// Pause and resume the variant.
variant.filter().set_revision_enabled(variant_revision, false).await?;
variant.filter().set_revision_enabled(variant_revision, true).await?;

// Another definition on a running group is refused.
match desk.filter().configure(ConsumerFilter::json(transitions_only)).await {
    Err(error) if error.filter_reason() == Some(FilterErrorReason::Conflict) => {}
    other => panic!("expected conflict, got {other:?}"),
}

// Release: the desk receives every record again.
let released = desk.filter().release().await?;
binding = await desk.filter().get()

# Draft on the desk's own filter. The desk keeps running revision 1.
draft = await desk.filter().revise(binding["revision"], ls.ConsumerFilter.json(transitions_only))
revisions = await desk.filter().revisions(page=0, page_size=10)

# Run the variant in its own group, with its own offsets.
variant = topic.consumer_group("anomaly-desk-transitions")
created = await variant.create(filter=ls.ConsumerFilter.json(transitions_only))

# Pause and resume the variant.
await variant.filter().set_revision_enabled(created["filter"]["revision"], False)
await variant.filter().set_revision_enabled(created["filter"]["revision"], True)

# Another definition on a running group is refused.
try:
    await desk.filter().configure(ls.ConsumerFilter.json(transitions_only))
except ls.FilterError as error:
    assert error.reason == "conflict"

# Release: the desk receives every record again.
released = await desk.filter().release()

Every catalog change carries an operation ID. A retry by the same caller with the same ID and mutation returns the first retained outcome. Reusing an ID with another mutation or caller is rejected. configure_as(operation_id, policy) in Rust (configureAs in TypeScript, configure(..., operation_id=) in Python) sends a configuration under an ID you recorded first. Results stay in durable, indexed plane storage, even after control-topic messages expire, so retry protection survives a recreated control topic. Recent results stay in memory up to the cache limit, and a miss reads durable storage.

Each change is authorized again when the catalog applies it, against the catalog as it stands at that point of its log, so a scoped deny holds even on a node that had not seen the group yet.

A group ID is scoped to its topic incarnation. Deleting and recreating a group, a topic, or a stream makes a new identity that inherits no policy. A reader of the old ID does not join the replacement.

Errors

Filter failures carry a typed reason. Rust exposes it through error.filter_reason(), Python through FilterError.reason, TypeScript through FilterExecutionError.reason.

ReasonWhat it means
conflictA precondition failed. The group runs another definition, the expected revision moved, or an acknowledgment names a policy the group no longer holds
source_changedThe source history or the group's identity changed. The reader restarts that partition from its stored offset once and raises this error so you know
revision_disabledThe active revision is paused. New server reads stop until it is enabled. Already delivered records can still be acknowledged
catalog_unavailableThe plane cannot prove its catalog is current, or the required configuration has not reached this node yet. Retryable
membership_staleThe consumer no longer owns the partition. The SDK refreshes its assignment
not_primaryThe node is no longer the partition primary. The SDK resolves the route again
not_foundThe group does not exist, or a verb that needs a policy met an unbound group
forbiddenThe caller lacks the filter grant on the group path or read permission on the source
unsupportedThe server does not serve this operation, such as the catalog without laser-plane, or a filter its evaluator version or codec does not match
capacity_exhaustedA configured catalog resource limit is reached, such as definitions, revisions, group identities, or stored revision bytes
version_skewThe client and the server speak different filter op versions

A stop fault or a record larger than the page cap blocks its partition. The reader delivers every earlier match first, then raises a filter fault with the partition and offset.

A managed server that serves filters but not group-aware reads is refused when you build a group consumer, with an unsupported error that asks for a server upgrade. A capability probe that failed is a retryable error. Neither falls back to a native poll, because that would read past the group's policy.

Key operations

VerbWhat it does
topic.consumer_group(name) / consumer_group_id(id)The group handle, by name or native ID. Free to construct
group.create().filter(filter).build()Create the group, optionally with its filter. Idempotent
group.info()ID, name, exact identity, and the active binding
group.consumer()The normal consumer. Runs the group's policy, batch_length bounds page and scan
group.reader()The match-oriented reader. count, max_examined, start, partition, read_mode, local_guard, idle_interval, max_unacked_pages
next_page() / next_record()Read the next page or record
ack(record) / ack_through(record) / ack_page(page)Mark records done, which stores progress
group.filter().configure(filter)Give an existing group its policy. configure_with takes one of its own revisions, configure_as an operation ID
group.filter().get()The active binding, or none
group.filter().revisions(page, size) / revise(expected, filter)List revisions, draft a new one
group.filter().set_revision_enabled(revision, enabled)Pause or resume without changing content
group.filter().release()Remove the policy. Consumers then receive every record
group.filter().delete()Delete the group's own filter with every revision, releasing the group first
group.filter().test(payload, headers) / preview(partition)Judge a sample or stored records without storing state

TypeScript spells these consumerGroup, consumerGroupId, nextPage, ackThrough, configureWith, setRevisionEnabled, and so on. Python passes options as keyword arguments and returns catalog replies as dicts.

A local guard also checks records marked unevaluated. The server includes its decoder bounds so the guard can reproduce a size or depth fault. A marker without the applicable pass policy is rejected. Older servers without these bounds cannot verify pass-through decode faults locally.

Permissions

Grant resources for a group's filter use the group path stream/topic/group. filter:read inspects and lists, filter:write drafts a revision, filter:admin configures, pauses, resumes and releases. Native source-read permission still controls consumption, and a consumer does not need to inspect the definition to run its group's policy.

In the LaserData Cloud Console, open a topic and its consumer group. The group's filter is shown above its members, with its revision, digest and policy generation. You can configure an unbound group there with the same visual and schema-aware editor, pause or resume the active revision, release the policy, and preview it. Unavailable policy information is shown as unavailable, never as "no filter".

Limits

LimitValue
Encoded filter size8 KiB
Sample header value1 KiB
Expression nodes128
Nesting depth8
Items in one in list64
Text pattern1 KiB
Compiled glob or regex programs per filter4
Budget per compiled program256 KiB
Records in one page1000
Unacknowledged record-bearing pages per partition1024 by default, configurable per reader
Bytes in one page8 MiB
Records one preview examines1,000 by default, configurable up to 10,000
Records one preview returns100
Live saved definitions1,024
Revisions per definition256
Retained group identities4,096
Stored revision bytes, all filters together64 MiB
Dropped filters kept as tombstones, oldest deleted past the limit4,096
Results cached in memory by default256
Concurrent mutations per plane by default32
Concurrent mutations per caller by default8

Catalog storage, memory use, and request admission are configured separately. The definition, tombstone, revision, and group-identity values above are defaults. Their LD_PLANE_FILTER_MAX_DEFINITIONS, LD_PLANE_FILTER_MAX_TOMBSTONES, LD_PLANE_FILTER_MAX_REVISIONS, LD_PLANE_FILTER_MAX_GROUP_IDENTITIES, and LD_PLANE_FILTER_MAX_CATALOG_BYTES settings accept 0 for no limit. LD_PLANE_FILTER_OUTCOME_CACHE_ENTRIES=0 disables the cache while keeping durable retries.

LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS and LD_PLANE_FILTER_MAX_CONCURRENT_MUTATIONS_PER_ACTOR set positive admission limits. When busy, the plane returns a retryable refusal. These limits count requests currently in progress, not historical operations.

The system control topic defaults to no expiry. Set LD_PLANE_CONTROL_EXPIRY_SECS to 604800 for seven days, 2592000 for thirty days, or 0 for no expiry. Startup applies the configured value to existing topics. With finite topic retention, recovery uses the durable plane snapshot and remaining log. Expiring topic messages does not delete group filters, bindings, or operation results.

LD_PLANE_DLQ_EXPIRY_SECS and LD_PLANE_CHANGES_EXPIRY_SECS keep the plane's dead-letter and change topics for one day by default.

IGGY_PLANE_FILTERS_ENABLED turns the filter evaluator on for a deployment. Group reads of unbound groups work without it, a bound group is refused until it is on. IGGY_PLANE_FILTERS_MAX_CONCURRENT allows 16 concurrent scans per shard by default, and a scan past it gets a retryable refusal. IGGY_PLANE_FILTERS_MAX_EXAMINED_BYTES bounds the bytes one page scan may examine at 64 MiB. IGGY_PLANE_FILTERS_PREVIEW_MAX_EXAMINED sets how many records a preview may examine, 1,000 by default with a ceiling of 10,000.

IGGY_PLANE_FILTERS_MAX_ROUND_RECORDS caps one internal owner read at 1,000 records by default, the same ceiling an ordinary poll has, with a configurable range of 1 to 1,000. Fewer owner reads per page cost less server CPU: on a one-million-record trial the cap of 1,000 used about 19% less server CPU than 128 with the same results. After the first read of a page, the next read shrinks to what the remaining examined-byte budget covers at the average record size seen so far. The first owner read probes at most 128 records. A sudden increase in frame size can still exceed the byte budget by one owner read. The scan yields the server every 1 MiB of records and sizes the next read by how selective the filter was so far, so a sparse filter can scan further within the page budget. A request's max_examined caps the budget below the server's. Consumer-group reads also respect the remaining desired match count to preserve the native polling fence. IGGY_PLANE_FILTERS_MAX_PAGE_ROUNDS bounds owner reads per page (256 by default). IGGY_PLANE_FILTERS_MAX_ROUNDS separately bounds previews. A page also stops at its configured round and byte budgets, so a 1,000-record request can return fewer matches. Byte accounting includes record frames and is checked between owner reads, so a page can pass it by one read. Previews size their owner reads the same way and also check it before each record. These bounds limit work per page. They don't establish a throughput or latency guarantee.

Group policies reuse cached definitions. Each broker shard caches resolved policies and watches its own plane for catalog changes. Before reusing an entry, it confirms the current catalog version through the sidecar. A delayed watch cannot keep an old binding active after a completed local catalog change. Changes made on another node first reach the local plane through control-log replay. A plane restart invalidates prior version tokens. The server also checks current reader grants before reuse. An unbound answer is cached for at most IGGY_PLANE_FILTERS_UNBOUND_CONFIRM_MS.

IGGY_PLANE_FILTERS_POLICY_CONFIRM_TIMEOUT_MS bounds each version confirmation at 1,000 ms by default. IGGY_PLANE_FILTERS_POLICY_CACHE_ENTRIES bounds entries per shard, with a default of 1,024 and a maximum of 65,536. IGGY_PLANE_FILTERS_POLICY_CACHE_TTL_MS sets the maximum entry lifetime, with a default of 60 seconds and a maximum of 10 minutes. Set either to zero to resolve every request. Entries hold weak references to compiled filters. Cache hits avoid repeated definition transfer and decoding, while a small version confirmation still runs for each hit. Watch failures and confirmation failures disable reuse. The cache confirmation counter separates those requests from full catalog lookups.

Payload formats and future codecs

Filter fields inside JSON, CBOR, Avro, and Protobuf without downloading every record. All formats use the same predicates, fault policies, previews, and acknowledgment rules. The server returns the original payload bytes, so consumers keep their existing decoders.

FormatFilter constructorValue rules
JSONConsumerFilter.json(expression)Exact 64-bit integers, finite numbers, arrays, objects, strings, booleans, and null
CBORConsumerFilter.cbor(expression)Definite-length values, text map keys, byte strings as integer arrays. No tags, duplicate keys, or non-finite numbers
AvroConsumerFilter.avro(expression, schema_ids)Registered writer schemas. Logical types use their physical values, such as integer timestamps or decimal bytes
ProtobufConsumerFilter.protobuf(expression, schema_ids)Registered descriptor sets and message names. Field names follow the schema, enums are numbers, and bytes are integer arrays. A field with presence tracking (proto2, optional, messages) is absent when unset. A proto3 scalar or repeated field without it always exists, at its default when unset, so counter == 0 matches
Headers onlyConsumerFilter.headers_only(expression)Typed user headers, with no payload decoding

The table uses Python constructor names. Rust uses ConsumerFilter::avro(expression, schema_ids) and TypeScript uses ConsumerFilter.avro(expression, schemaIds). The CDC examples show complete typed producers and group readers in all three languages.

Register Avro and Protobuf schemas once, then reference their IDs. Put the allowed writer IDs in the filter's schema_refs. Stamp each published record with its writer ID through schema_id in Rust/Python or schemaId in TypeScript. This sets the existing agdx.sid header. A filter can allow several immutable writer schemas of the same codec. The schema IDs participate in its digest. A missing, unlisted, or incompatible schema ID follows the filter's foreign_policy when the predicate needs the payload: the record is skipped by default, or delivered unevaluated under pass.

The broker resolves schemas through the existing plane schema catalog before scanning records. Rust, Python, and TypeScript local guards use the same references. CBOR needs no registered schema. Protobuf groups are outside this profile. Empty repeated and map fields are present with their default empty values. Use explicit timestamp coercions for Avro integer timestamps.

Byte and nesting limits bound decoding. The server announces its evaluator version and supported codecs. SDKs reject unsupported capabilities explicitly. Header-only routing works with opaque payloads in any format. Additional codecs can reuse the decoder boundary and predicate evaluator without changing Iggy transport or storage. Each new codec must define value semantics and pass the shared conformance corpus.

Availability and cost

Group reads, tests, and previews run on the LaserData Iggy fork, in Laser Stack and LaserData Cloud. Configuring a group's filter needs laser-plane. With the plane disabled, a group reads unfiltered only when no retained control topic or required catalog position indicates prior catalog history. Otherwise, the read returns catalog_unavailable. create with a filter is refused before the group is created when catalog support is unavailable. Check capabilities().filters before you use them: catalog for configuration, group_policy_reads for group-aware consumers. Original Apache Iggy remains supported for native streaming.

Measure bandwidth savings and server cost together. The broker still scans source records per group. Header-only policies avoid decoder work. Payload policies add decoding and evaluation while reducing bytes sent to consumers. Use the Frostline benchmark documentation to measure CPU, memory, sample counts and p50/p99/p99.9 latency on your workload.

The example figures at the top describe recorded workloads. They are not a general latency guarantee or a replacement for benchmarks on your deployment.

In TypeScript, acknowledge the original reader-owned record or page. Copying its payload to another task does not transfer acknowledgment ownership.

On this page