Streams & Tasks
Gnok provides two complementary building blocks for data pipelines: streams expose table changes (change data capture), and tasks execute SQL on schedules or when a condition is true. Use streams when a consumer needs a managed Iceberg snapshot watermark; use tasks for scheduling and dependency chains.
Streams
A stream is a change-tracking object that exposes inserts, updates, and deletes on a table since its watermark. Streams use Iceberg snapshot history and do not copy or modify the source table.
How Streams Work
When you create a stream on a table, Gnok records the current Iceberg snapshot as the starting watermark. As INSERT, UPDATE, and DELETE operations commit new snapshots, the stream exposes the delta between its watermark and a fixed current snapshot. An ordinary SELECT previews that delta without changing the watermark. CONSUME STREAM returns the delta and advances the watermark after successful query execution.
Each row returned by a stream includes metadata columns:
| Column | Type | Description |
|---|---|---|
METADATA$ACTION | VARCHAR | The type of change: INSERT, DELETE |
METADATA$ISUPDATE | BOOLEAN | TRUE if the row is part of an UPDATE (appears as DELETE + INSERT pair) |
METADATA$ROW_ID | VARCHAR | Unique identifier for the row |
Stream Types
| Type | Tracks | Use Case |
|---|---|---|
| Standard | Inserts, updates, and deletes | Auditing and incremental consumers |
| Append-only | Inserts only | Streaming ingestion, event logs, append-only tables |
Append-only streams are more efficient because they track new data files and do not calculate row-level deletes. Use them only when the source is operationally append-only.
Creating Streams
-- Standard stream (tracks all changes)
CREATE STREAM order_changes ON TABLE sales.orders;
-- Append-only stream (tracks inserts only)
CREATE STREAM event_log_stream ON TABLE analytics.events
APPEND_ONLY = TRUE;
-- Include the table's current rows in the first pending batch
CREATE STREAM order_bootstrap ON TABLE sales.orders
SHOW_INITIAL_ROWS = TRUE;
The stream must be created in the same catalog and schema as its source table. Creating or reading a stream also requires read access to the source table. By default, rows that existed before CREATE STREAM are the baseline and are not returned; use SHOW_INITIAL_ROWS = TRUE when the first batch must contain them.
Consuming Streams
Use SELECT to preview pending changes. Previewing is repeatable and never advances the watermark.
-- Preview pending changes without advancing the watermark
SELECT * FROM order_changes;
-- Return pending changes, then advance the watermark
CONSUME STREAM order_changes;
CONSUME STREAM plans a bounded delta to a fixed snapshot, executes it, and then advances the watermark with compare-and-swap. If another consumer advances the same stream concurrently, one operation fails instead of silently overwriting the watermark.
The current release does not make INSERT ... SELECT FROM <stream> or MERGE ... USING <stream> atomic with stream advancement. A Gnok watermark also cannot share a transaction with an unrelated external sink. Use idempotent destination writes and reconciliation, or use explicit READ CHANGES ... BETWEEN SNAPSHOT ... ranges with an application-managed checkpoint when an external consumer requires replayable boundaries.
Checking for Data
Use SYSTEM$STREAM_HAS_DATA to check whether the source has a relevant snapshot after the stream watermark. This is useful for polling and as a task condition:
SELECT SYSTEM$STREAM_HAS_DATA('order_changes');
-- Returns TRUE or FALSE
Managing Stream Offsets
-- Acknowledge and discard the pending batch without returning it
ALTER STREAM order_changes ADVANCE;
ALTER STREAM ... ADVANCE is intentionally destructive: it moves the watermark to the source table's current snapshot and discards the pending batch. Use it only as an explicit skip or recovery action.
Listing and Dropping Streams
-- List all streams in the current schema
SHOW STREAMS;
-- List streams on a specific table
SHOW STREAMS ON TABLE sales.orders;
-- Drop a stream
DROP STREAM order_changes;
Metadata and Operational Limits
Use the metadata columns to interpret each change:
- An insert is one row with
METADATA$ACTION = 'INSERT'. - A delete is one row with
METADATA$ACTION = 'DELETE'. - With valid Iceberg identifier fields, an update is a
DELETE/INSERTpair withMETADATA$ISUPDATE = TRUEand the same stableMETADATA$ROW_ID. - Without identifier fields, Gnok preserves insert/delete bag semantics, but cannot distinguish an update pair from independent changes.
Plan for these current limits:
- One query can read only one independent stream.
- A query cannot reference a stream and its source table directly in the same statement because their snapshot boundaries differ.
- Stream reads depend on the source snapshots between the watermark and current head. Coordinate snapshot expiration with the oldest live stream watermark.
- A lost client response after
CONSUME STREAMis ambiguous: the watermark might have committed even though the client did not receive the response. - Give each independent consumer its own stream. Compare-and-swap protects a shared stream from silent concurrent advancement, but it does not turn one watermark into a broadcast queue.
Tasks
A task is a scheduled SQL statement that executes automatically at defined intervals or in response to data availability. Tasks support dependency chains, enabling multi-step pipelines where child tasks run after their parents complete.
How Tasks Work
Tasks are defined with a SQL body and a schedule. When the schedule triggers, Gnok checks any optional conditions (such as whether a stream has data) and, if satisfied, executes the SQL statement using the specified warehouse. Tasks are created in a suspended state and must be explicitly resumed before they will run.
Schedules
Tasks support two scheduling modes:
Interval scheduling -- run at a fixed frequency:
CREATE TASK refresh_summary
WAREHOUSE = etl_wh
SCHEDULE = '5 MINUTES'
AS
INSERT INTO analytics.hourly_summary
SELECT date_trunc('hour', event_time) AS hour, COUNT(*) AS event_count
FROM analytics.events
WHERE event_time > (SELECT MAX(hour) FROM analytics.hourly_summary)
GROUP BY 1;
CRON scheduling -- run on a cron expression:
CREATE TASK daily_report
WAREHOUSE = reporting_wh
SCHEDULE = 'USING CRON 0 9 * * MON-FRI UTC'
AS
CALL generate_daily_report();
Stream Triggers
Tasks can be conditioned on stream data availability with a WHEN SYSTEM$STREAM_HAS_DATA(...) clause. The condition controls whether the task body runs; it does not consume the stream or advance its watermark.
Do not use a task to imply atomic MERGE-and-acknowledge semantics. In the current release, target DML and stream advancement are separate operations. If a scheduled workflow needs both, design it for idempotent replay and reconciliation, or use explicit snapshot ranges with its own durable checkpoint.
Task Dependencies
Child tasks run after their parent task completes successfully. Use the AFTER clause to define dependencies:
-- Parent task: extract raw data
CREATE TASK extract_raw
WAREHOUSE = etl_wh
SCHEDULE = '10 MINUTES'
AS
INSERT INTO staging.raw_events (event_id, event_type, payload, event_time)
SELECT event_id, event_type, payload, event_time
FROM raw.events
WHERE event_time > COALESCE(
(SELECT MAX(event_time) FROM staging.raw_events),
TIMESTAMP '1970-01-01 00:00:00'
);
-- Child task: transform (runs after extract_raw completes)
CREATE TASK transform_events
WAREHOUSE = etl_wh
AFTER extract_raw
AS
INSERT INTO analytics.processed_events
SELECT event_id, event_type, parse_json(payload) AS parsed
FROM staging.raw_events
WHERE processed = FALSE;
-- Grandchild task: load aggregates (runs after transform_events completes)
CREATE TASK load_aggregates
WAREHOUSE = etl_wh
AFTER transform_events
AS
MERGE INTO analytics.event_counts AS target
USING (
SELECT event_type, COUNT(*) AS cnt
FROM analytics.processed_events
WHERE event_time > CURRENT_DATE
GROUP BY event_type
) AS source
ON target.event_type = source.event_type
WHEN MATCHED THEN UPDATE SET target.cnt = source.cnt
WHEN NOT MATCHED THEN INSERT (event_type, cnt) VALUES (source.event_type, source.cnt);
Task Lifecycle
Tasks are created in a SUSPENDED state. You must explicitly resume a task before it will execute:
-- Resume a task (starts scheduling)
ALTER TASK extract_raw RESUME;
-- Suspend a task (stops scheduling)
ALTER TASK extract_raw SUSPEND;
When resuming a task tree (parent + children), resume children first, then the parent. When suspending, suspend the parent first, then children.
Manual Execution
Trigger a task immediately without waiting for the next scheduled run:
EXECUTE TASK extract_raw;
This executes the task once, regardless of schedule or WHEN conditions.
Listing and Dropping Tasks
-- List all tasks in the current schema
SHOW TASKS;
-- View task execution history
SELECT * FROM TABLE(INFORMATION_SCHEMA.TASK_HISTORY())
ORDER BY scheduled_time DESC
LIMIT 20;
-- Drop a task
DROP TASK load_aggregates;
Common Patterns
CDC Inspection and Acknowledgement
Create a stream, inspect its pending changes as often as needed, and explicitly acknowledge a batch:
CREATE STREAM customer_cdc ON TABLE raw.customers;
SELECT SYSTEM$STREAM_HAS_DATA('customer_cdc');
SELECT
customer_id,
name,
email,
METADATA$ACTION,
METADATA$ISUPDATE,
METADATA$ROW_ID
FROM customer_cdc;
CONSUME STREAM customer_cdc;
See Build an Incremental Inventory Change Feed with Gnok CDC Streams for an end-to-end example and production delivery guidance.
ELT Pipeline
A scheduled task chain that extracts, transforms, and loads data:
-- Extract: runs every 10 minutes and uses an application timestamp watermark
CREATE TASK elt_extract
WAREHOUSE = etl_wh
SCHEDULE = '10 MINUTES'
AS
INSERT INTO staging.extracted (id, name, amount, exchange_rate, updated_at)
SELECT id, name, amount, exchange_rate, updated_at
FROM source.records
WHERE updated_at > COALESCE(
(SELECT MAX(updated_at) FROM staging.extracted),
TIMESTAMP '1970-01-01 00:00:00'
);
-- Transform: runs after extract completes
CREATE TASK elt_transform
WAREHOUSE = etl_wh
AFTER elt_extract
AS
INSERT INTO staging.transformed
SELECT id, UPPER(name) AS name, amount * exchange_rate AS amount_usd
FROM staging.extracted
WHERE NOT processed;
-- Load: runs after transform completes
CREATE TASK elt_load
WAREHOUSE = etl_wh
AFTER elt_transform
AS
MERGE INTO production.summary USING staging.transformed AS src
ON production.summary.id = src.id
WHEN MATCHED THEN UPDATE SET name = src.name, amount_usd = src.amount_usd
WHEN NOT MATCHED THEN INSERT (id, name, amount_usd) VALUES (src.id, src.name, src.amount_usd);
-- Resume in order: children first, then parent
ALTER TASK elt_load RESUME;
ALTER TASK elt_transform RESUME;
ALTER TASK elt_extract RESUME;
Alerting
A task that checks a condition and records alerts:
CREATE TASK error_rate_alert
WAREHOUSE = monitoring_wh
SCHEDULE = '5 MINUTES'
AS
INSERT INTO ops.alerts (alert_type, message, created_at)
SELECT
'HIGH_ERROR_RATE',
'Error rate exceeded 5% in the last 10 minutes: ' || CAST(error_rate AS VARCHAR),
CURRENT_TIMESTAMP
FROM (
SELECT
COUNT_IF(status = 'ERROR') * 100.0 / COUNT(*) AS error_rate
FROM api.request_log
WHERE request_time > DATEADD('MINUTE', -10, CURRENT_TIMESTAMP)
)
WHERE error_rate > 5.0;
ALTER TASK error_rate_alert RESUME;