LaserData Cloud
Laser SDKAdvanced

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

  1. Register a projection. It names the fields to extract from each payload.
  2. Apply a binding. It connects a (stream, topic) source to the projection and names the operational index.
  3. Publish to the topic as usual. The managed plane extracts the fields from each new record and writes a row to the index.
  4. 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 under name. field_at("cpu_pct", "/cpu/value") indexes cpu.value as cpu_pct.
  • field_typed(name, type) and field_at_typed(name, pointer, type) add a storage type hint for typed-column backends. Rust uses FieldType::Int, Float, Bool, or Text. 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 is Any, 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 is name.vN, for example readings_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. IndexSchemaBuilder builds one with field, field_at, vector_field, and inline_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 with for_topic, for_topics, name_contains, id_prefix, and search (camelCase in TypeScript) and finish with fetch(). Python passes topic=, topics=, name_contains=, id_prefix=, and search= to list(..). 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:

CallRustTypeScriptPython
build()PanicsThrows InvalidErrorRaises InvalidError
try_build() / tryBuild()Returns Err(InvalidError)Returns undefinedRaises InvalidError

A hand-written binding uses the fields source, allowed_projections, default_projection, backend, index, notify, and retention (camelCase in TypeScript).

Retention policies

PolicyRustPythonTypeScript
Mirror the log, the defaultRetentionPolicy::MirrorLog{"kind": "mirror_log"}{ kind: "mirrorLog" }
Keep foreverRetentionPolicy::Keep{"kind": "keep"}{ kind: "keep" }
Keep until the source topic is deletedRetentionPolicy::KeepUntilSourceDeleted{"kind": "keep_until_source_deleted"}{ kind: "keepUntilSourceDeleted" }
Time to live after materializationRetentionPolicy::TimeToLive { ttl_micros }{"kind": "time_to_live", "ttl_micros": n}{ kind: "timeToLive", ttlMicros: n }
Newest rows onlyRetentionPolicy::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 passes name= and version= to register(..). 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, and filter_prefix compare any indexed field. Chained filters combine with AND.
  • filter(expr) adds a composed condition. Build it with Filter::pred, Filter::all, Filter::any, and Filter::negate in Rust, ls.Filter.pred, all, any, and negate in Python, and filterPred, filterAll, filterAny, and filterNegate in TypeScript.
  • message_type(value) matches the reserved message_type field.
  • time_range(start, end) matches the reserved ts field, in epoch microseconds. The range includes both ends, and start must be before end.

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) and order_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 and has_more is 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, and count_distinct work on every backend. stddev and percentile(field, fraction) need a columnar backend.
  • Each aggregate writes to a result field named after itself, such as count or avg. TypeScript aggregate methods take an optional alias as their last argument.
  • agg_as adds 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 in window_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}")
  • nearest(embedding, top_k) runs a nearest-neighbor search on the embedding field. Rust and TypeScript use nearest_in(field, embedding, top_k) (nearestIn) for another field. Python passes field= to nearest.
  • 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's score. Lexical search needs the keyword capability, 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 needs select_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.

LanguageCallDialect
Rustraw_sql(dialect, sql), or raw_sql_with(dialect, sql, params) with positional parametersRequired, a SqlDialect
TypeScriptrawSql(sql, params?, dialect?)Defaults to "data_fusion"
Pythonraw_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 a ConversationId. 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:

LanguageAccessors
Rust and Pythonresult.value(row, name), value_text, value_u64, value_i64, and field_index(name)
TypeScriptqueryResultValue(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_typed returns one page of decoded values. Rust decodes JSON into your type with fetch_typed::<T>() and takes another decoder with fetch_typed_with::<C, T>(). Python decodes JSON into Python values, and fetch_typed_with(codec) takes any object with decode(data). TypeScript takes a codec, as in fetchTyped(new Json<Reading>()).
  • fetch_one returns at most one decoded value. Rust and Python also have fetch_one_with. TypeScript fetchOne(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 += 1
let 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 += 1

A 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

SituationWhat you get
The deployment does not serve queriesUnsupported error, before sending
An unadvertised consistency level, lexical search, cursor paging, status, or cancelUnsupported 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 aggregateInvalid error, before sending
rows() or rows_typed() without max_rowsInvalid error
Read-your-writes that cannot catch up in timeStale query error
A typed read of a row without a payloadConfig error
The index does not exist yetAn error until the managed plane creates it
register with a graph projectionInvalid error

Key operations

VerbWhat it does
Projection::builder(id)Start a projection (TypeScript and Python: ProjectionBuilder(id))
.field / .fields / .field_atIndex top-level fields or a JSON pointer
.field_typed / .field_at_typedIndex 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 / dropManage the projection registry
ProjectionBinding::builder()Start a binding (TypeScript and Python: ProjectionBindingBuilder())
.source / .allow / .default_projectionWhich topic, which projections, and the projection for untagged records
.index / .backend / .retention / .notifyIndex name, backend, row lifetime, and change records
bindings().apply / removeApply or remove a binding
schemas().register / get / list / dropManage 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_targetOpen a query
where_eq and the filter_* familyMatch and filter
order_asc / order_desc / limit / offset / cursorSort and page
count, sum, avg, group_by, window, havingAggregate
nearest / textVector and lexical search
read_your_writes / consistencyChoose read consistency
fetch / fetch_typed / fetch_one / fetch_allRun the query
max_rows(n).rows() / rows_typed()Walk rows across pages under a ceiling
status / cancelRead or stop the execution
execute_query / query_page / query_status / cancel_queryLower-level query lifecycle
raw_sqlOne read-only SQL statement on a SQL backend

On this page