Changes in depth
The change feed reader, change record fields, resume rules, delivery guarantees, and capability checks
This is the full reference for Changes. Start there for the short version.
The change feed tells you when an operational view advanced. Each change record says which index moved, over which source offsets, and how many rows landed. It does not carry the rows. Read the feed, then query the index only when it moved. The feed is an ordinary topic read by offset over the connection you already have, so it adds no new connection or push channel.
Turn it on
Call notify() on the projection binding. After each committed projector batch for that binding, the managed plane publishes one change record. A binding without notify() serves queries only. Watch the same index name that you query.
The feed needs the watch capability, which Laser Stack and LaserData Cloud advertise. Reading the rows afterward also needs the query capability. In Rust, the reader needs the watch feature, which the managed feature includes.
Where records go
Change records go to the changes topic on the ops stream, _agdx, the stream the managed services use for their own records. On a deployment that reports stream_tenancy, each stream has its own change topic, stream:<stream>/_agdx/changes. A client with a default stream reads its stream's topic, and each record names its stream. A client with a default stream also scopes the index name it filters on, the same way a query scopes it.
Open a reader
| Language | Open a reader for one index | Read every index | Returns |
|---|---|---|---|
| Rust | laser.watch().index(name).records()? | laser.watch().records()? | Result<WatchReader, LaserError>, synchronously |
| TypeScript | await laser.watch().index(name).records() | await laser.watch().records() | Promise<WatchReader> |
| Python | laser.watch(index=name) | laser.watch() | The WatchReader directly, not awaited |
When the deployment does not publish the feed, opening the reader fails with an unsupported error: records() in Rust and TypeScript, and watch() in Python. The SDK never waits on a topic that nothing writes to.
.index(name) keeps one index. The filter runs in the client, so the reader still reads every record on the topic. The reader skips any record that does not decode as a change record.
Read records
Your application owns the loop. The SDK tracks the position between polls and has no push callbacks.
- Call
poll(). It returns every matching record that arrived since the last poll, or an empty list when the reader is caught up. - Handle the records. Query the index if it moved.
- If nothing arrived, wait briefly and poll again.
To read one record at a time, Rust's stream() returns a futures::Stream, TypeScript's stream() returns an async generator, and the Python reader is itself an async iterator. All three end when the reader is caught up. In TypeScript and Python a later loop on the same reader resumes where the last one stopped. Rust's stream() takes ownership of the reader, so save its offsets first if you need them.
const feed = await laser.watch().index("readings_v1").records()
for await (const change of feed.stream()) {
console.log(`${change.index} advanced ${change.rows} row(s) on partition ${change.partitionId}`)
}use futures::StreamExt;
let feed = laser.watch().index("readings_v1").records()?;
let mut changes = std::pin::pin!(feed.stream());
while let Some(change) = changes.next().await {
let change = change?;
println!(
"{} advanced {} row(s) on partition {}",
change.index, change.rows, change.partition_id
);
}feed = laser.watch(index="readings_v1")
async for change in feed:
print(f"{change.index} advanced {change.rows} row(s) on partition {change.partition_id}")Change record fields
| Field | TypeScript | Meaning |
|---|---|---|
index | index | The index that advanced |
partition_id | partitionId | The source partition |
from_offset | fromOffset | First source offset in the batch, inclusive |
to_offset | toOffset | Last source offset in the batch, inclusive |
rows | rows | Rows the batch wrote |
stream | stream | The source stream, set only when the deployment keeps one feed per stream |
v | v | The change record version, currently 1 |
In TypeScript, fromOffset and toOffset are bigint values and stream is optional. In Rust, stream is an Option<String>. In Python, it is None when absent.
Resume after a restart
Save the reader's offsets and restore them in a new reader. Where you store them is up to you.
| Language | Read the position | Restore it |
|---|---|---|
| Rust | feed.offsets(), a slice with one u64 per partition | .from_offsets(saved) on a new reader, with a Vec<u64> |
| TypeScript | The feed.offsets property, a ReadonlyMap<number, bigint> from partition to offset | .fromOffsets(saved) on a new reader |
| Python | The feed.offsets property, a list of integers | .from_offsets(saved), or laser.watch(index=name, from_offsets=saved) |
Each offset is the next one to read on that partition. A new reader without restored offsets starts at the beginning of the feed. In Rust, a saved list shorter than the partition count reads the missing partitions from the start.
const saved = feed.offsets
// After a restart
const resumed = (await laser.watch().index("readings_v1").records()).fromOffsets(saved)let saved = feed.offsets().to_vec();
// After a restart
let mut resumed = laser
.watch()
.index("readings_v1")
.records()?
.from_offsets(saved);saved = feed.offsets
# After a restart
resumed = laser.watch(index="readings_v1", from_offsets=saved)Delivery guarantees
- Notifications are best effort after commit. A lost notification does not lose projected data, because the rows are already in the view.
- A reader that falls behind the feed's retention window misses the records that expired. Query the view directly to catch up.
- The feed reports progress. It does not prove that a specific write is in the view. To make sure a query includes your own write, use a read-your-writes query. See Consistency.
- The feed reports operational view changes only. It does not report key-value or session changes. Sessions have their own change feed on the session reads. Destination checkpoints have their own lifecycle.
Key operations
| Verb | What it does |
|---|---|
notify() on a binding | Turn on change records for that index |
watch().index(name).records() | Open a reader for one index (Rust, TypeScript) |
watch(index=name) | Open a reader for one index (Python) |
watch().records() / watch() | Read every index on one feed |
poll() | Read the records that arrived since the last poll |
stream() | Read records one at a time until caught up (Rust, TypeScript. Python: async for) |
offsets / from_offsets(saved) | Save and restore the position across restarts (TypeScript: fromOffsets) |