Changes
Poll view progress and refresh queries when new data is available
Changes reports when an operational view advances. Poll the feed, then query the index when new rows are available. Both operations use the existing connection.
Built for
Use Changes for live interfaces, cache invalidation, and reactive pipelines.
How it works
Enable notifications on a projection binding to publish a change record after each materialized batch. Without that flag, the binding serves queries only. Watch the same index name that you query.
Rust and TypeScript use watch().index(name).records(). Python uses watch(index=name).
The application owns the polling loop:
- Call
.poll(). - Process the records received since the last poll.
- If no records arrive, wait briefly before polling again.
The SDK tracks position between polls. It does not provide push callbacks. Rust's .stream() returns a futures::Stream and ends when caught up. Python's async reader follows the same pattern.
Save offsets() to preserve per-partition progress across restarts. Restore it through .from_offsets(saved) on the next reader. Without restored offsets, a watcher starts at the feed's beginning.
Each record describes a committed range: index, source partition, offsets, and row count. to_offset reports the new watermark, the position through which processing advanced. Query the index to retrieve the resulting rows.
This feed reports operational view changes. It does not report general key-value or run-registry changes. Destination checkpoints use their own lifecycle.
The reader skips malformed records. .index(name) excludes other indexes. Use watch().records() to read all indexes from one feed. Query and watch capabilities are both required, and the example tests them before reading.
Quick example
// orders_v1 declared with notify: true above this block
const feed = await laser.watch().index("orders_v1").records()
while ((await feed.poll()).length > 0) {
// drain every batch a previous run left behind
}
await laser
.topic("orders")
.publish()
.json({ id: 4, status: "paid", total: 20 })
.send()
let changes = await feed.poll()
while (changes.length === 0) {
await new Promise((resolve) => setTimeout(resolve, 200))
changes = await feed.poll()
}
for (const change of changes) {
console.log(
`view advanced: ${change.rows} row(s), ` +
`source offsets ${change.fromOffset}..${change.toOffset}`
)
}let mut feed = laser.watch().index("orders_v1").records()?;
while !feed.poll().await?.is_empty() {
// drain every batch a previous run left behind
}
laser.topic("orders").publish().json(&order)?.send().await?;
loop {
let changes = feed.poll().await?;
if !changes.is_empty() {
for change in changes {
println!(
"view advanced: {} row(s), source offsets {}..{}",
change.rows, change.from_offset, change.to_offset
);
}
break;
}
tokio::time::sleep(Duration::from_millis(200)).await;
}feed = laser.watch(index="orders_v1")
while await feed.poll():
pass # drain every batch a previous run left behind
await laser.topic("orders").publish(
{"id": 4, "total": 20, "status": "paid"}
).send()
while not (changes := await feed.poll()):
await asyncio.sleep(0.2)
for change in changes:
print(
f"view advanced: {change.rows} row(s), "
f"source offsets {change.from_offset}..{change.to_offset}"
)Complete examples: Rust, Python, and TypeScript.
Key operations
| Verb | What it does |
|---|---|
| Bind a projection with notify on | Turns on change delivery for that index |
watch().index(name).records() | Open a change reader (Rust, TypeScript) |
watch(index=name) | Open a change reader, one call (Python) |
watch().records() | Watch every index on one feed |
.poll() | Drain whatever landed since the last poll |
.offsets() / .from_offsets(saved) | Persist and restore the feed position across restarts |
.stream() | The reader as a futures::Stream, ending once caught up (Rust) |
| A change record's row and offset fields | How much changed and where, not the changed rows themselves |
Running it
Changes requires laser-plane, available in Laser Stack and LaserData Cloud.