Build an Incremental Inventory Change Feed with Gnok CDC Streams
An inventory table may contain millions of products, but a downstream stock monitor usually needs only the rows that changed. Repeatedly scanning the full table wastes compute and makes deletes especially difficult to detect.
Gnok CDC streams turn the Iceberg snapshots behind a table into a bounded change feed. In this tutorial, we will create an inventory stream, inspect inserts, updates, and deletes, and acknowledge a batch after processing it.
This is outbound table CDC: the source data already lives in Gnok. For inbound data, see Streaming Ingestion. For historical snapshot ranges without a managed watermark, see Snapshots, Time Travel, and Rollback.
The Use Case
Suppose a retailer stores its current warehouse stock in an Iceberg table. Inventory writes continue normally, while an operations application periodically needs to answer:
- Which products were added?
- Which stock levels changed?
- Which products were removed?
- Has anything changed since the previous acknowledged batch?
The stream stores a watermark, not a second copy of the table. Gnok derives changes between that watermark and a fixed current Iceberg snapshot when the stream is queried.
Create the Inventory Table
Create a catalog, schema, and table for the example:
CREATE CATALOG retail;
CREATE SCHEMA retail.inventory;
CREATE TABLE retail.inventory.stock_levels (
sku VARCHAR PRIMARY KEY,
warehouse VARCHAR,
on_hand BIGINT,
updated_at TIMESTAMP
);
The primary key becomes an Iceberg identifier field. That gives CDC a stable logical row identity, which allows Gnok to recognize an update as a related DELETE and INSERT pair.
Seed the current state before creating the stream:
INSERT INTO retail.inventory.stock_levels
(sku, warehouse, on_hand, updated_at)
VALUES
('SKU-100', 'east', 12, TIMESTAMP '2026-08-27 09:00:00'),
('SKU-200', 'west', 5, TIMESTAMP '2026-08-27 09:00:00');
Start Tracking Changes
Create a standard stream in the same catalog and schema as its source table:
CREATE STREAM retail.inventory.stock_changes
ON TABLE retail.inventory.stock_levels;
By default, the current table snapshot becomes the stream's starting watermark. The two seed rows therefore establish the baseline and do not appear as new changes.
Now simulate normal inventory activity:
INSERT INTO retail.inventory.stock_levels
(sku, warehouse, on_hand, updated_at)
VALUES
('SKU-300', 'east', 20, TIMESTAMP '2026-08-27 09:05:00');
UPDATE retail.inventory.stock_levels
SET on_hand = 9,
updated_at = TIMESTAMP '2026-08-27 09:06:00'
WHERE sku = 'SKU-100';
DELETE FROM retail.inventory.stock_levels
WHERE sku = 'SKU-200';
Each committed write creates a new Iceberg snapshot. The source table continues to represent current inventory; the stream represents the changes since its watermark.
Check and Preview the Pending Batch
An inexpensive availability check can prevent polling clients from running an empty change query:
SELECT SYSTEM$STREAM_HAS_DATA(
'retail.inventory.stock_changes'
);
It returns one has_data value, which is now true.
Preview the pending changes with an ordinary SELECT:
SELECT
sku,
warehouse,
on_hand,
METADATA$ACTION,
METADATA$ISUPDATE,
METADATA$ROW_ID
FROM retail.inventory.stock_changes
ORDER BY sku, METADATA$ACTION;
The result has this shape:
sku | warehouse | on_hand | METADATA$ACTION | METADATA$ISUPDATE | METADATA$ROW_ID |
|---|---|---|---|---|---|
| SKU-100 | east | 12 | DELETE | true | same ID as next row |
| SKU-100 | east | 9 | INSERT | true | same ID as previous row |
| SKU-200 | west | 5 | DELETE | false | stable row ID |
| SKU-300 | east | 20 | INSERT | false | stable row ID |
Iceberg represents an update as removal of the old row plus insertion of the new row. METADATA$ISUPDATE = true and the shared row ID let a consumer pair those two records. A standalone delete or insert has METADATA$ISUPDATE = false.
Run the same SELECT again and it returns the same pending batch. Ordinary stream queries never advance the watermark, so they are safe for inspection, validation, and debugging.
Consume the Batch
After the consumer is ready to acknowledge this bounded batch, use the explicit consuming statement:
CONSUME STREAM retail.inventory.stock_changes;
CONSUME STREAM returns the pending rows and then advances the watermark with a compare-and-swap operation. The watermark advances only after the change query succeeds. A concurrent consumer cannot silently overwrite a different watermark; one of the operations fails and must be retried.
Confirm that the stream is current:
SELECT SYSTEM$STREAM_HAS_DATA(
'retail.inventory.stock_changes'
);
SELECT *
FROM retail.inventory.stock_changes;
The availability check now returns false, and the preview returns no rows. Future table commits create a new pending batch.
Choose the Right Delivery Pattern
The stream watermark is transactional inside Gnok, but it cannot participate in the transaction of an unrelated external database, search index, or message broker. That boundary matters when designing a production consumer.
- Use ordinary
SELECTfor repeatable previews and validation. It never acknowledges data. - Use
CONSUME STREAMwhen returning the change result is itself the consumption step, or when the surrounding application can reconcile its destination. - Make downstream operations idempotent where possible.
METADATA$ROW_ID,METADATA$ACTION, and the business key are useful deduplication inputs. - Treat a lost client response as ambiguous: the server may have returned the rows and committed the watermark even if the client did not receive the final response.
- For an external sink that requires replayable, application-managed checkpoints, use explicit
READ CHANGES ... BETWEEN SNAPSHOT ...ranges and persist the completed snapshot boundary in that application.
Gnok does not currently make a separate INSERT ... SELECT FROM <stream> or MERGE ... USING <stream> atomic with stream advancement. Do not describe that two-system workflow as exactly-once without an additional transactional handoff design.
Production Checklist
Before using a stream in a production workflow:
- Grant least privilege. The principal needs
CREATE STREAMon the schema and read access to the source table. Stream reads re-check source-table authorization. - Use identifier fields. A primary key or valid Iceberg identifier fields provide stable row IDs and update pairing. Without them, changes remain valid insert/delete records, but updates cannot be distinguished from unrelated deletes and inserts.
- Keep the stream beside its source. The stream and table must currently share a catalog and schema.
- Retain required snapshots. A stream derives its delta from Iceberg snapshot history. Coordinate snapshot expiration with the oldest live stream watermark.
- Monitor the watermark. Use
SHOW STREAMS ON TABLE retail.inventory.stock_levelsandSYSTEM$STREAM_HAS_DATAto observe lag and pending work. - Advance deliberately.
ALTER STREAM ... ADVANCEdiscards the currently pending batch by moving the watermark without returning rows. Reserve it for an intentional skip or recovery procedure. - Avoid competing consumers. Compare-and-swap prevents silent watermark corruption, but independent consumers should generally have independent streams.
Clean Up
When the example is no longer needed:
DROP STREAM retail.inventory.stock_changes;
DROP TABLE retail.inventory.stock_levels;
For syntax, stream variants, operational limits, and task scheduling, see the Streams & Tasks guide.