Queries and views in depth
Every projection, binding, and query option, with defaults, limits, errors, and language differences
This is the full reference for Queries and views. Start there for the short version.
A view is a queryable table that the managed plane keeps up to date from messages in the log. A projection says which fields to extract. A binding says which topic feeds the projection and which index the rows land in. Queries read the index. Views need Laser Stack or LaserData Cloud. On standalone Apache Iggy the query capability is off and every call on this page returns an unsupported error.
How a view is built
- Register a projection. It names the fields to extract from each payload.
- Apply a binding. It connects a
(stream, topic)source to the projection and names the operational index. - Publish to the topic as usual. The managed plane extracts the fields from each new record and writes a row to the index.
- Query the index.
Registration and binding both publish control commands. The managed plane applies them asynchronously, so the calls return before the index exists. Register the projection before you apply a binding that names it, or the binding is rejected. A query fails until the index exists, then returns an empty page.
A client with a default stream scopes projection IDs, index names, and the index a query names to that stream, as stream:<stream>/<name>. Listings return your stream's projections without the prefix. Schema registration and lookup use the stream's own registry. See Stream-scoped resource names.
In Rust, views need the projections and query features. The managed feature includes both.
Register a projection
All three SDKs build a projection with the same builder: Projection::builder(id) in Rust, new ProjectionBuilder(id) in TypeScript, and ls.ProjectionBuilder(id) in Python. TypeScript uses camelCase method names. build() returns the projection, which is a dict in Python. Pass it to projections().register(projection).
field(name)indexes a top-level JSON field under the same name.fields([...])indexes several top-level fields.field_at(name, pointer)indexes the value at an RFC 6901 JSON pointer undername.field_at("cpu_pct", "/cpu/value")indexescpu.valueascpu_pct.field_typed(name, type)andfield_at_typed(name, pointer, type)add a storage type hint for typed-column backends. Rust usesFieldType::Int,Float,Bool, orText. Python and TypeScript use"int","float","bool", or"text". The embedded engine ignores the hint.vector_field(pointer)extracts an embedding vector for nearest-neighbor search. Without it, the vector is read from/embedding.content_type(..)selects the decoder, such as JSON or Avro. The builder default isAny, a best-effort decode.inline_payload()keeps a copy of the original payload with each row, so typed reads decode the row without reading the log. This is the builder default.index_only()keeps only the declared fields.version(n)sets the schema version. The default is 1. The recommended ID form isname.vN, for examplereadings_v1.v1.name(..)sets the display name. It is separate from the ID, so you can rename a projection without changing the ID that producers send.extraction(schema)replaces the whole extraction plan.IndexSchemaBuilderbuilds one withfield,field_at,vector_field, andinline_payload.
A hand-written object or dict with the same fields also works. It states the payload choice with inline_payload (inlinePayload) on the extraction and inline_payload_default (inlinePayloadDefault) on the projection. A dict that omits them stores no payload, unlike the builder.
A producer can also ask for one record's payload to be kept with its row with inline_payload() on the publish request (TypeScript: inlinePayload()).
Limits:
- A record can carry at most 32
.index(..)values. The SDK refuses more when you send. - The managed plane decodes and copies payloads up to 8 MiB. A larger record keeps only the values its
.index(..)headers set. The projection's fields are not extracted from it, and the row does not copy the payload. - A record with no indexed fields produces no row.
The log keeps the original bytes according to the topic's retention. Rows have their own retention, set on the binding, so a row can outlive the message that produced it. Only declared fields are queryable. Other fields stay in the inline payload, when it is kept, or in the log.
projections().register(..) rejects a graph projection with an invalid error. Register those with register_graph as shown in Graph.
Read and change the registry
projections().get(id)returns one projection with its bindings, or nothing when the ID is unknown.projections().list()browses the registry. Rust and TypeScript narrow it withfor_topic,for_topics,name_contains,id_prefix, andsearch(camelCase in TypeScript) and finish withfetch(). Python passestopic=,topics=,name_contains=,id_prefix=, andsearch=tolist(..). No filter lists every projection.projections().drop(id)stops materialization for that projection. Rows already written stay.
Writes are applied asynchronously. Poll get(id) to see when a registration has been applied.
Bind a topic
A binding selects one (stream, topic) source and the projections allowed to read it. The builder is ProjectionBinding::builder() in Rust, new ProjectionBindingBuilder() in TypeScript, and ls.ProjectionBindingBuilder() in Python. Pass the result to bindings().apply(binding).
source(stream, topic)selects the topic. A binding has exactly one source.selector(source)takes a prebuilt source instead.allow(projection_id)adds a permitted projection. Call it once per projection.default_projection(id)handles records that name no projection. Without a default, those records are skipped.index(name)names the operational index. Without it, the index takes the source topic name.backend(binding)pins the index to an exact backend resource ID and generation. Without it, the index uses the replica-local backend.retention(policy)sets how long rows live, independent of topic retention. Without it, the binding uses the deployment default, which mirrors the log unless an operator changed it.notify()turns on change records for this binding. Notifications are off by default.
A producer tags a record for one projection with projection_ref(id) on the publish request (TypeScript: projectionRef). It sets the agdx.ref header. A record tagged for a projection outside the allowed set is skipped or dead-lettered, according to the deployment's policy. The header keeps your local projection ID, and the managed plane resolves it against the bindings of the record's own stream, so a record never selects another stream's projection.
The build step differs by language when no source is set:
| Call | Rust | TypeScript | Python |
|---|---|---|---|
build() | Panics | Throws InvalidError | Raises InvalidError |
try_build() / tryBuild() | Returns Err(InvalidError) | Returns undefined | Raises InvalidError |
A hand-written binding uses the fields source, allowed_projections, default_projection, backend, index, notify, and retention (camelCase in TypeScript).
Retention policies
| Policy | Rust | Python | TypeScript |
|---|---|---|---|
| Mirror the log, the default | RetentionPolicy::MirrorLog | {"kind": "mirror_log"} | { kind: "mirrorLog" } |
| Keep forever | RetentionPolicy::Keep | {"kind": "keep"} | { kind: "keep" } |
| Keep until the source topic is deleted | RetentionPolicy::KeepUntilSourceDeleted | {"kind": "keep_until_source_deleted"} | { kind: "keepUntilSourceDeleted" } |
| Time to live after materialization | RetentionPolicy::TimeToLive { ttl_micros } | {"kind": "time_to_live", "ttl_micros": n} | { kind: "timeToLive", ttlMicros: n } |
| Newest rows only | RetentionPolicy::MaxRows { rows } | {"kind": "max_rows", "rows": n} | { kind: "maxRows", rows: n } |
In TypeScript, n is a bigint. Mirror the log removes rows when Iggy removes the messages that produced them. Mirror the log and keep until the source topic is deleted both drop the rows when the source topic is deleted. Keep, time to live, and newest rows only keep their rows after the topic is gone. A policy kind that a newer server sends and your SDK does not know decodes as Unknown. Do not apply a binding you read back with an unknown policy, because its original kind is lost.
The index and the topic have separate names. The readings topic can feed the readings_v1 index, so you can version the view without renaming the topic that producers use.
bindings().remove(source, projection_ref) stops materialization for that source. Rows already written stay. Pass a projection ID to stop routing into that one projection only. Rust takes Option<String> for the second argument. TypeScript and Python make it optional.
Register and bind, in code
import { ContentType, ProjectionBindingBuilder, ProjectionBuilder } from "@laserdata/laser-sdk"
const projection = new ProjectionBuilder("readings_v1.v1")
.name("readings_v1")
.version(1)
.contentType(ContentType.Json)
.indexOnly()
.fields(["host", "cpu", "status"])
.build()
await laser.projections().register(projection)
const binding = new ProjectionBindingBuilder()
.source("telemetry", "readings")
.allow("readings_v1.v1")
.defaultProjection("readings_v1.v1")
.index("readings_v1")
.retention({ kind: "timeToLive", ttlMicros: 7n * 24n * 3_600_000_000n })
.notify()
.build()
await laser.bindings().apply(binding)use laser_sdk::prelude::full::*;
let projection = Projection::builder("readings_v1.v1")
.name("readings_v1")
.version(1)
.content_type(ContentType::Json)
.index_only()
.fields(["host", "cpu", "status"])
.build();
laser.projections().register(projection).await?;
let binding = ProjectionBinding::builder()
.source("telemetry", "readings")
.allow("readings_v1.v1")
.default_projection("readings_v1.v1")
.index("readings_v1")
.retention(RetentionPolicy::TimeToLive {
ttl_micros: 7 * 24 * 3_600_000_000,
})
.notify()
.build();
laser.bindings().apply(binding).await?;import laser_sdk as ls
projection = (
ls.ProjectionBuilder("readings_v1.v1")
.name("readings_v1")
.version(1)
.content_type("json")
.index_only()
.fields(["host", "cpu", "status"])
.build()
)
await laser.projections().register(projection)
binding = (
ls.ProjectionBindingBuilder()
.source("telemetry", "readings")
.allow("readings_v1.v1")
.default_projection("readings_v1.v1")
.index("readings_v1")
.retention({"kind": "time_to_live", "ttl_micros": 7 * 24 * 3_600_000_000})
.notify()
.build()
)
await laser.bindings().apply(binding)After the binding is applied, every new message on telemetry/readings updates the view.
Two ways to index a field
With a registered projection, the producer sends a normal payload and the managed plane extracts the fields through the projection's pointers. The producer needs no projection details.
The producer can also set an indexed value directly. .index(key, value) on a publish request adds an agdx.idx.<key> header with a string value. Use it for raw payloads, for writer schemas the plane cannot decode, or for a constant on every record in a batch. A header wins over an extracted value with the same field name. A record can carry at most 32 of these headers.
Writer schemas
laser.schemas() is the writer schema registry for Avro, Protobuf, and JSON Schema payloads. Registration is synchronous: the managed plane checks that the definition compiles, allocates the next free ID, and returns it. IDs are permanent and concurrent callers never get the same one. A definition that does not compile returns an unsupported error and allocates nothing.
- Rust and TypeScript call
register(source), optionally chain.name(..)and.version(..), and finish with.send(). Python passesname=andversion=toregister(..). Name and version are stored labels only. get(id)returns the schema at that ID, active or dropped, or nothing when the ID is free.list()returns every known schema, active and dropped.drop(id)marks the schema dropped. Records stamped with that ID keep decoding, and the ID can never be reused for a different definition. Dropping an unknown ID does nothing.
Producers stamp the ID with .schema_id(id) on the publish request (TypeScript: schemaId).
Wait for the view
Materialization is asynchronous. A message you just published can be missing from the next query.
- Wait for a query to succeed after you apply a binding. The examples do this before they publish, so new records reach a running projector.
- Use the change feed to learn when an index advanced, then query.
- Use a read-your-writes query when one query must see your own earlier writes. See Consistency.
Query an index
laser.query(index) opens a query against an operational index. Chain conditions, then run it with a terminal such as fetch(). TypeScript uses camelCase names, such as whereEq, filterGte, and countDistinct.
Match and filter
where_eq(field, value)matches an indexed field exactly. It is the cheap path that index keys answer directly.filter_eq,filter_ne,filter_gt,filter_gte,filter_lt,filter_lte,filter_in,filter_contains, andfilter_prefixcompare any indexed field. Chained filters combine with AND.filter(expr)adds a composed condition. Build it withFilter::pred,Filter::all,Filter::any, andFilter::negatein Rust,ls.Filter.pred,all,any, andnegatein Python, andfilterPred,filterAll,filterAny, andfilterNegatein TypeScript.message_type(value)matches the reservedmessage_typefield.time_range(start, end)matches the reservedtsfield, in epoch microseconds. The range includes both ends, andstartmust be beforeend.
Value types differ. Rust takes anything that converts into a TypedValue, such as &str, i32, i64, or f64. Python takes plain Python values. TypeScript whereEq takes a string or a TypedValue, and the filter* comparisons take a TypedValue object such as { kind: "long", value: 80n }. filterContains and filterPrefix take a string.
Sort and page
order_asc(field)andorder_desc(field)sort. Several calls sort by several fields in order.limit(n)sets the page size. The default is 50 and the maximum is 1,000. Zero or a value above 1,000 fails with an invalid error before anything is sent.offset(n)skips rows.cursor(c)resumes from a cursor that a previous page returned. Setting one clears the other.with_total()asks for an exact total count. It costs a full scan on a wide filter. Without it, the total is absent andhas_moreis still exact.
Each result has a page with offset, limit, total, has_more, and next_cursor (hasMore and nextCursor in TypeScript). next_cursor is present exactly when has_more is true. Aggregate and vector queries always return one page.
Aggregate
count,sum,avg,min,max, andcount_distinctwork on every backend.stddevandpercentile(field, fraction)need a columnar backend.- Each aggregate writes to a result field named after itself, such as
countoravg. TypeScript aggregate methods take an optional alias as their last argument. agg_asadds an aggregate with a chosen alias, so you can return several of the same kind. Rust:agg_as(func, field, fraction, alias). Python:agg_as(func, alias, field=, fraction=). TypeScript:aggAs(func, alias, { field, fraction }).group_by([...])groups the aggregate.window(field, every_micros)buckets it into tumbling windows over a timestamp field, and each result row carries the bucket start inwindow_start.having(expr)filters aggregate output. Its fields name an aggregate alias or a group key, not raw row fields. It needs an aggregate.- An aggregate query cannot also select fields or ask for payloads.
import { queryResultValueI64, queryResultValueText } from "@laserdata/laser-sdk"
const byStatus = await laser.query("readings_v1").count().groupBy(["status"]).fetch()
for (const row of byStatus.rows) {
const status = queryResultValueText(byStatus, row, "status")
const count = queryResultValueI64(byStatus, row, "count")
console.log(`${status ?? "?"}: ${String(count ?? 0n)}`)
}let by_status = laser
.query("readings_v1")
.count()
.group_by(["status"])
.fetch()
.await?;
for row in &by_status.rows {
let status = by_status.value_text(row, "status").unwrap_or_default();
let count = by_status.value_i64(row, "count").unwrap_or(0);
println!("{status}: {count}");
}by_status = await laser.query("readings_v1").count().group_by(["status"]).fetch()
for row in by_status.rows:
status = by_status.value_text(row, "status")
count = by_status.value_i64(row, "count") or 0
print(f"{status}: {count}")Search
nearest(embedding, top_k)runs a nearest-neighbor search on theembeddingfield. Rust and TypeScript usenearest_in(field, embedding, top_k)(nearestIn) for another field. Python passesfield=tonearest.text(query)runs a lexical search over every text-hinted indexed field.text_in(field, query)narrows it to one field. Relevance lands in each row'sscore. Lexical search needs thekeywordcapability, and the SDK refuses it before sending when the deployment does not advertise it.
Select fields and payloads
select_fields([...])returns only some fields.distinct()returns distinct rows over the selected fields. It needsselect_fields.with_payload()returns the stored payload bytes on each row.
Raw SQL
raw_sql runs one read-only SQL statement on the index's backend. It works on SQL backends only, is not portable, and cannot be combined with any structured condition, sort, aggregate, selection, or payload request.
| Language | Call | Dialect |
|---|---|---|
| Rust | raw_sql(dialect, sql), or raw_sql_with(dialect, sql, params) with positional parameters | Required, a SqlDialect |
| TypeScript | rawSql(sql, params?, dialect?) | Defaults to "data_fusion" |
| Python | raw_sql(sql, params=None, dialect="data_fusion") | Defaults to "data_fusion" |
The dialects are DataFusion, Postgres, MySQL, and SQLite.
The statement reads the queried index by the same name the query uses. A stream-scoped name such as stream:orders/orders.v1 holds : and /, so quote it: FROM "stream:orders/orders.v1", or with backticks on MySQL.
Consistency
Queries are eventually consistent by default. read_your_writes() waits until the projector reaches the end of the source log as it was when the query started. If the projector does not catch up in time, the query fails with a stale error instead of returning older data. consistency(level) takes eventual, read-your-writes, or strong. Strong adds agreement across replicas.
Read-your-writes and strong need backend support. The SDK checks capabilities.query.consistency and returns an unsupported error before sending when the level is not served. Check for a stale error with error.is_stale() in Rust, isStale(error) in TypeScript, and the stale attribute in Python.
Scope a query
conversation(id)returns only the rows one conversation produced. The deployment indexes every message's conversation header, so this works on any projection without producer changes. Rust takes aConversationId. TypeScript and Python take a string.fork(id)reads a fork's copy-on-write view instead of the trunk. See Key-value state.
Deadlines
Each query has one deadline for all its pages, 30 seconds after the request was created by default. The server never extends it. Set a relative deadline with deadline(duration) in Rust and deadline(timeout_ms) in milliseconds in Python and TypeScript. deadline_micros(..) (TypeScript: deadlineMicros) sets an absolute time in epoch microseconds in all three SDKs.
Read results
A result has fields, rows, page, and context. Each row's values follow the order of fields. Read values by field name with the accessors:
| Language | Accessors |
|---|---|
| Rust and Python | result.value(row, name), value_text, value_u64, value_i64, and field_index(name) |
| TypeScript | queryResultValue(result, row, name), queryResultValueText, queryResultValueU64, queryResultValueI64, and queryResultFieldIndex(result, name) |
Use typed accessors for logic and the text accessor for display. value_u64 and value_i64 return nothing when the value is missing, out of range, or not an integer. TypeScript typedValueDiagnosticText(value) and Python ls.typed_value_diagnostic_text(value) render any typed value as text. See Managed data for the type rules.
Typed reads
Typed reads decode the stored payload of each row, so they need a view that keeps payloads. A row without a payload fails.
fetch_typedreturns one page of decoded values. Rust decodes JSON into your type withfetch_typed::<T>()and takes another decoder withfetch_typed_with::<C, T>(). Python decodes JSON into Python values, andfetch_typed_with(codec)takes any object withdecode(data). TypeScript takes a codec, as infetchTyped(new Json<Reading>()).fetch_onereturns at most one decoded value. Rust and Python also havefetch_one_with. TypeScriptfetchOne(codec)takes the codec directly.fetch_all()reads every matching row into memory.fetch_all_typed()(TypeScript:fetchAllTyped(codec)) decodes them. Use these only when the result fits in memory.
Walk many rows
rows() walks matching rows across pages. It needs an explicit ceiling, max_rows(n) (TypeScript: maxRows(n)), and fails with an invalid error without one. It stops at the ceiling or the last page, whichever comes first. Rust returns a walker that you advance with next().await. TypeScript returns an async iterable. Python returns an asynchronous iterator. rows_typed() does the same with decoded payloads (TypeScript: rowsTyped(codec)).
let degraded = 0
const walk = laser.query("readings_v1").whereEq("status", "degraded").maxRows(5_000).rows()
for await (const _row of walk) degraded += 1let mut rows = laser
.query("readings_v1")
.where_eq("status", "degraded")
.max_rows(5_000)
.rows()?;
let mut degraded = 0;
while rows.next().await?.is_some() {
degraded += 1;
}rows = laser.query("readings_v1").where_eq("status", "degraded").max_rows(5_000).rows()
degraded = 0
async for _row in rows:
degraded += 1A zero ceiling returns no rows and sends no request in all three SDKs. TypeScript rejects a negative or unsafe integer ceiling. The walk fetches pages of limit rows, 50 by default.
Execution identity, status, and cancel
Every query carries an execution ID. Rust reads it with execution_id(), TypeScript with executionId(), and Python from the execution_id property. status() reads the execution state and cancel() asks the server to stop it.
Cursor paging, status, and cancellation each have their own capability flag: cursor_paging, execution_status, and cancellation on capabilities.query (camelCase in TypeScript). The SDK returns an unsupported error before sending when the deployment does not advertise the one you call.
The lower-level calls on the client are execute_query(query), query_page(execution_id, cursor, deadline_micros), query_status(execution_id), and cancel_query(execution_id) (TypeScript: executeQuery, queryPage, queryStatus, cancelQuery). Each checks that the reply carries the execution ID you asked for, and a reply for another execution fails with a protocol error. When the server advertises a different query protocol version, the SDK fails with a typed version error before sending.
Build a query value
To build a whole query instead of chaining on laser.query(..), use Query::builder() in Rust, new QueryBuilder() in TypeScript, or ls.QueryBuilder() in Python, then pass the result to execute_query (TypeScript: executeQuery). The execution ID, target, and absolute deadline in epoch microseconds are required, and build() fails with an invalid error when one is missing. The builder has one setter per query field. In all three SDKs, into_query() (TypeScript: intoQuery()) turns a chained request into the same value, and query_target(target) (TypeScript: queryTarget) starts a chained request from an explicit target.
Operational views and destinations
A binding keeps an operational index on the deployment's backend. A materialization destination is a separate resource with its own identity, source scope, generation, checkpoints, and lifecycle. It is not another target of the binding.
A query names its target. laser.query(index) reads an operational index. laser.query_lakehouse(destination, generation) (TypeScript: queryLakehouse) reads one destination generation. On a lakehouse target, at_snapshot(id) and at_timestamp_micros(ts) (TypeScript: atSnapshot, atTimestampMicros) read an earlier table state. Both values must be positive. On an operational target they fail with an invalid error: Rust returns it, TypeScript throws it, and Python raises it. A lakehouse query cannot read a fork. What a backend can serve depends on what it advertises. See Managed data.
Columnar data enters the log as Arrow IPC. Publish one complete stream per message with arrow_ipc(bytes, metadata) on a publish request (TypeScript: arrowIpc). The SDK checks the metadata and that the payload length matches it before sending. The managed plane enforces the rest of the acceptance policy.
Errors
| Situation | What you get |
|---|---|
| The deployment does not serve queries | Unsupported error, before sending |
| An unadvertised consistency level, lexical search, cursor paging, status, or cancel | Unsupported error naming the feature, before sending |
limit of zero or above 1,000, an empty time range, both offset and cursor, or having without an aggregate | Invalid error, before sending |
rows() or rows_typed() without max_rows | Invalid error |
| Read-your-writes that cannot catch up in time | Stale query error |
| A typed read of a row without a payload | Config error |
| The index does not exist yet | An error until the managed plane creates it |
register with a graph projection | Invalid error |
Key operations
| Verb | What it does |
|---|---|
Projection::builder(id) | Start a projection (TypeScript and Python: ProjectionBuilder(id)) |
.field / .fields / .field_at | Index top-level fields or a JSON pointer |
.field_typed / .field_at_typed | Index a field with a storage type hint |
.vector_field(pointer) | Extract an embedding vector |
.inline_payload() / .index_only() | Keep the payload with each row, or only the indexed fields |
projections().register / get / list / drop | Manage the projection registry |
ProjectionBinding::builder() | Start a binding (TypeScript and Python: ProjectionBindingBuilder()) |
.source / .allow / .default_projection | Which topic, which projections, and the projection for untagged records |
.index / .backend / .retention / .notify | Index name, backend, row lifetime, and change records |
bindings().apply / remove | Apply or remove a binding |
schemas().register / get / list / drop | Manage writer schemas |
publish().index(key, value) | Set an indexed value at publish time |
publish().projection_ref(id) | Tag a record for one allowed projection |
query(index) / query_lakehouse / query_target | Open a query |
where_eq and the filter_* family | Match and filter |
order_asc / order_desc / limit / offset / cursor | Sort and page |
count, sum, avg, group_by, window, having | Aggregate |
nearest / text | Vector and lexical search |
read_your_writes / consistency | Choose read consistency |
fetch / fetch_typed / fetch_one / fetch_all | Run the query |
max_rows(n).rows() / rows_typed() | Walk rows across pages under a ceiling |
status / cancel | Read or stop the execution |
execute_query / query_page / query_status / cancel_query | Lower-level query lifecycle |
raw_sql | One read-only SQL statement on a SQL backend |