Skip to main content

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:

ColumnTypeDescription
METADATA$ACTIONVARCHARThe type of change: INSERT, DELETE
METADATA$ISUPDATEBOOLEANTRUE if the row is part of an UPDATE (appears as DELETE + INSERT pair)
METADATA$ROW_IDVARCHARUnique identifier for the row

Stream Types​

TypeTracksUse Case
StandardInserts, updates, and deletesAuditing and incremental consumers
Append-onlyInserts onlyStreaming 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.

Delivery boundary

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/INSERT pair with METADATA$ISUPDATE = TRUE and the same stable METADATA$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 STREAM is 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;