Snowflake courseLesson 6 of 12
Snowflake course · Lesson 6 of 12
Snowflake Streams and Tasks: Change Data Capture and Scheduling
Use Snowflake streams to capture inserts, updates and deletes, and tasks to process them on a schedule or trigger: offsets, staleness, task graphs and error handling.
On this page
- Stream object basics
- Standard versus append-only streams
- Insert-only streams
- Stream offsets
- CDC with streams
- A runnable analogue in PostgreSQL
- Stale streams
- Task fundamentals
- Scheduled tasks with CRON
- Task graphs (DAGs)
- Serverless tasks
- The task plus stream pattern
- Playbook: a stream and task CDC pipeline
- Task error handling
- Practice questions
- Key takeaways
A stream records what changed in a table since you last looked; a task runs SQL on a schedule or when data arrives. Together they are Snowflake’s built-in way to process changes incrementally: load raw data, capture new and changed rows with a stream, and apply them to modelled tables with a task, without an external orchestrator. They come up constantly in Data Engineer interviews because they test whether you understand change data capture (CDC), exactly-once processing and failure handling.
All Snowflake SQL in this lesson is written from the documentation and was not executed. One section includes a runnable PostgreSQL analogue that imitates a stream, clearly labelled.
Stream object basics
A stream is an object you create on a source object, most often a table:
-- Snowflake SQL (not executed here)
CREATE OR REPLACE STREAM customers_stream ON TABLE raw.customers;
It does not copy data. It stores an offset: a point in the source table’s version history. When you query the stream, Snowflake compares the table’s current version with the version at the offset and returns the rows that changed, with three extra metadata columns:
| Column | Meaning |
|---|---|
METADATA$ACTION |
INSERT or DELETE |
METADATA$ISUPDATE |
TRUE when the row is half of an update: an UPDATE appears as a DELETE of the old row and an INSERT of the new one, both with TRUE |
METADATA$ROW_ID |
A unique, immutable ID for the source row, so you can track it across changes |
Streams can be created on standard tables, views (including secure views), dynamic tables, directory tables, event tables, external tables and Iceberg tables, with type restrictions described below. Creating the first stream on a table enables change tracking on it, which adds hidden columns that record row identity; on a view, change tracking must be enabled on the underlying tables.
Because a stream relies on the table’s history, it can only see changes still inside the table’s data retention period (Time Travel). That is the root of staleness, covered later.
Pitfalls
- Thinking a stream is a queue or a copy of the data. It is a pointer plus a comparison. If the source history it needs is gone, so are the changes.
- Expecting
SELECT * FROM streamto “use up” the changes. Querying alone never moves the offset.
In interviews
“What is a Snowflake stream?” A strong answer: a change-tracking object that stores an offset into a table’s version history and returns the rows changed since then, with METADATA$ACTION, METADATA$ISUPDATE and METADATA$ROW_ID; the offset advances only when the stream is consumed in a committed DML statement.
Standard versus append-only streams
| Standard (delta) | Append-only | |
|---|---|---|
| Captures | Inserts, updates and deletes | Inserts only |
| Supported on | Tables, views, dynamic tables, Iceberg tables managed by Snowflake, and others | The same standard sources |
| Result | Net changes between offset and now | Every row inserted since the offset |
| Cost to query | Higher: Snowflake joins inserted and deleted rows to compute the net change | Lower: only new rows are returned |
| Typical use | Keeping a modelled table in sync (dimensions, MERGE) | Append-only ingestion (events, logs) into staging |
Net changes matter for standard streams. Between two consumptions:
- a row inserted and then deleted does not appear at all;
- a row inserted and then updated appears as one
INSERTwith the final values (METADATA$ISUPDATE = FALSE); - a row updated three times appears as one
DELETEof the original values and oneINSERTof the final values, both withMETADATA$ISUPDATE = TRUE.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE STREAM events_stream ON TABLE raw.events APPEND_ONLY = TRUE;
An append-only stream ignores updates and deletes completely, including TRUNCATE. If rows already returned by an append-only stream are later deleted from the source, the stream does not tell you.
Pitfalls
- Using a standard stream on a high-volume, insert-only events table. You pay for delta computation you do not need.
- Using an append-only stream where updates matter. Late corrections in the source silently never reach the target.
In interviews
Explain net changes with the insert-then-delete example, and say you would choose append-only for raw event ingestion and standard for anything that must mirror updates and deletes.
Insert-only streams
Insert-only streams are the version for sources Snowflake does not manage, where it cannot see deletes reliably: external tables and externally managed Iceberg tables (and some other external formats). They track rows added by new files.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE STREAM ext_orders_stream
ON EXTERNAL TABLE lake.orders_ext
INSERT_ONLY = TRUE;
Behaviour to know:
- They capture new rows from files registered in the external table’s metadata (for example through
AUTO_REFRESHorALTER EXTERNAL TABLE ... REFRESH). - They do not record deletes. If a file is removed and another is added, the stream shows only the rows from the new file.
- If a file is overwritten in place, the rows may appear again as new inserts.
- Streams are not supported on partitioned external tables according to the current documentation; check yours.
Pitfalls
- Treating an insert-only stream as a complete CDC feed. It is a “new files” feed.
In interviews
Usually asked as “what stream type would you use on an external table?”. Answer: insert-only, the only type supported there, and explain the lack of delete tracking.
Stream offsets
The offset is the heart of a stream. Rules from the documentation:
- The offset advances only when the stream is used in a DML statement that commits.
INSERT ... SELECT FROM stream,MERGE ... USING stream,CREATE TABLE ... AS SELECT FROM streamandCOPY INTO <location>from a stream all count. A plainSELECTdoes not, even inside a transaction. - All rows are consumed at once. When the DML commits, the offset moves to the version at the start of the transaction, even if your DML filtered the stream with a
WHEREclause and only used some rows. Rows you filtered out are gone from the stream. - If the transaction rolls back, the offset does not move. Changes are processed again next time. This is what gives you exactly-once application when you consume a stream inside one DML statement or one transaction.
- Repeatable reads in a transaction. Inside an explicit transaction, every read of the stream sees the same set of changes, so you can use one stream in several statements (for example insert into two tables) before committing, and they consume the same changes.
- One stream per consumer. Two pipelines reading one stream compete: whoever consumes first advances the offset for both. Create a separate stream for each consumer.
-- Snowflake SQL (not executed here)
-- Consume one stream into two targets with the same set of changes
BEGIN;
INSERT INTO audit.customer_changes
SELECT *, CURRENT_TIMESTAMP() FROM customers_stream;
INSERT INTO staging.customers_latest
SELECT customer_id, name, tier FROM customers_stream
WHERE METADATA$ACTION = 'INSERT';
COMMIT; -- the offset advances once, here
Pitfalls
- Filtering a stream in the consuming statement and expecting the filtered-out rows to wait for next time.
- Testing a pipeline by running the
MERGEby hand. It consumes the changes the scheduled task was meant to process.
In interviews
The question “when does a stream offset move?” is a favourite. Say “on commit of a DML statement that reads the stream”, add the all-or-nothing consumption rule and the rollback guarantee, and mention one stream per consumer.
CDC with streams
A standard stream gives you everything needed to apply changes to a target with one MERGE:
- rows with
METADATA$ACTION = 'INSERT'are new or updated rows: insert or update the target; - rows with
METADATA$ACTION = 'DELETE'andMETADATA$ISUPDATE = FALSEare real deletes: delete from the target; - rows with
METADATA$ACTION = 'DELETE'andMETADATA$ISUPDATE = TRUEare the “before” image of an update: skip them, because the matchingINSERTcarries the new values.
-- Snowflake SQL (not executed here)
MERGE INTO analytics.dim_customer AS t
USING (
SELECT *
FROM customers_stream
WHERE NOT (METADATA$ACTION = 'DELETE' AND METADATA$ISUPDATE) -- drop update "before" images
) AS s
ON t.customer_id = s.customer_id
WHEN MATCHED AND s.METADATA$ACTION = 'DELETE' THEN DELETE
WHEN MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
UPDATE SET t.name = s.name, t.tier = s.tier, t.updated_at = CURRENT_TIMESTAMP()
WHEN NOT MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
INSERT (customer_id, name, tier, updated_at)
VALUES (s.customer_id, s.name, s.tier, CURRENT_TIMESTAMP());
Because a standard stream returns net changes, each key appears at most once as an INSERT and at most once as a DELETE, so the MERGE does not hit “duplicate row matched” errors from the stream itself. Duplicates in the source data (two rows with the same business key) are a separate problem: deduplicate them in the USING subquery, for example with QUALIFY ROW_NUMBER() OVER (PARTITION BY customer_id ORDER BY loaded_at DESC) = 1.
For an SCD Type 2 dimension, the same stream feeds a pattern that closes the current version and inserts a new one; see slowly changing dimensions.
A runnable analogue in PostgreSQL
Snowflake cannot run here. The following PostgreSQL 16 script imitates a standard stream so you can watch the mechanics: a trigger writes each change to a change table with action and is_update columns (like METADATA$ACTION and METADATA$ISUPDATE), and an offset table records the last change consumed. This is an analogy, not how Snowflake implements streams (Snowflake uses table versions, not triggers).
-- PostgreSQL analogue of a stream (runs on PostgreSQL 16)
CREATE TABLE src_customers (
customer_id INT PRIMARY KEY,
name TEXT,
tier TEXT
);
-- The "stream": a change log with Snowflake-style metadata columns
CREATE TABLE customers_changes (
change_id BIGSERIAL PRIMARY KEY,
action TEXT, -- like METADATA$ACTION: 'INSERT' or 'DELETE'
is_update BOOLEAN, -- like METADATA$ISUPDATE
customer_id INT,
name TEXT,
tier TEXT
);
-- The stream's offset: the last change already consumed
CREATE TABLE stream_offset (stream_name TEXT PRIMARY KEY, last_change_id BIGINT);
INSERT INTO stream_offset VALUES ('customers_stream', 0);
CREATE FUNCTION capture_change() RETURNS trigger AS $$
BEGIN
IF TG_OP IN ('DELETE', 'UPDATE') THEN
INSERT INTO customers_changes (action, is_update, customer_id, name, tier)
VALUES ('DELETE', TG_OP = 'UPDATE', OLD.customer_id, OLD.name, OLD.tier);
END IF;
IF TG_OP IN ('INSERT', 'UPDATE') THEN
INSERT INTO customers_changes (action, is_update, customer_id, name, tier)
VALUES ('INSERT', TG_OP = 'UPDATE', NEW.customer_id, NEW.name, NEW.tier);
END IF;
RETURN NULL;
END $$ LANGUAGE plpgsql;
CREATE TRIGGER customers_cdc AFTER INSERT OR UPDATE OR DELETE ON src_customers
FOR EACH ROW EXECUTE FUNCTION capture_change();
-- Target table kept in sync by the "task"
CREATE TABLE dim_customer (customer_id INT PRIMARY KEY, name TEXT, tier TEXT);
Make some changes and look at the raw change log since the offset:
INSERT INTO src_customers VALUES (1, 'Asha', 'gold'), (2, 'Ben', 'silver'), (3, 'Chen', 'bronze');
UPDATE src_customers SET tier = 'gold' WHERE customer_id = 2;
DELETE FROM src_customers WHERE customer_id = 3;
SELECT change_id, action, is_update, customer_id, tier
FROM customers_changes
WHERE change_id > (SELECT last_change_id FROM stream_offset WHERE stream_name = 'customers_stream')
ORDER BY change_id;
change_id | action | is_update | customer_id | tier
-----------+--------+-----------+-------------+--------
1 | INSERT | f | 1 | gold
2 | INSERT | f | 2 | silver
3 | INSERT | f | 3 | bronze
4 | DELETE | t | 2 | silver
5 | INSERT | t | 2 | gold
6 | DELETE | f | 3 | bronze
The update to customer 2 is a DELETE/INSERT pair with is_update true, exactly like a Snowflake stream row pair. A Snowflake standard stream would show the net result instead: customer 1 and customer 2 (with tier gold) as plain inserts, and nothing for customer 3, which was inserted and deleted in the same interval. The consuming step below computes that net change (the last change per key wins), applies it with MERGE, and advances the offset in the same transaction:
BEGIN;
WITH pending AS (
SELECT * FROM customers_changes
WHERE change_id > (SELECT last_change_id FROM stream_offset WHERE stream_name = 'customers_stream')
),
net AS (
SELECT DISTINCT ON (customer_id) customer_id, action, name, tier
FROM pending
ORDER BY customer_id, change_id DESC
)
MERGE INTO dim_customer AS t
USING net AS s
ON t.customer_id = s.customer_id
WHEN MATCHED AND s.action = 'DELETE' THEN DELETE
WHEN MATCHED AND s.action = 'INSERT' THEN UPDATE SET name = s.name, tier = s.tier
WHEN NOT MATCHED AND s.action = 'INSERT' THEN INSERT VALUES (s.customer_id, s.name, s.tier);
-- Advance the offset in the same transaction, as consuming a stream in DML does
UPDATE stream_offset
SET last_change_id = (SELECT COALESCE(MAX(change_id), 0) FROM customers_changes)
WHERE stream_name = 'customers_stream';
COMMIT;
SELECT * FROM dim_customer ORDER BY customer_id;
SELECT * FROM stream_offset;
customer_id | name | tier
-------------+------+------
1 | Asha | gold
2 | Ben | gold
stream_name | last_change_id
------------------+----------------
customers_stream | 6
If the transaction had failed, neither the target nor the offset would have changed, and the next run would process the same changes: the same guarantee a Snowflake stream gives. (In this simplified analogue a change committed by another session while the merge runs could be skipped; Snowflake avoids that by pinning the stream to the table version at the start of the transaction.)
Pitfalls
- Forgetting to filter out update “before” images, so updates delete the target row.
- Ignoring duplicate business keys in the source, which makes
MERGEfail or behave nondeterministically.
In interviews
Be ready to write this MERGE from memory and explain each WHEN clause. Interviewers also ask how you would make the pipeline exactly-once: consume the stream in one DML statement (or one transaction), so the offset and the target change together.
Stale streams
A stream becomes stale when its offset falls outside the source table’s data retention period. The history needed to compute the changes has been purged, so the unconsumed changes are lost and the stream can no longer be read.
How Snowflake helps:
- If a table’s
DATA_RETENTION_TIME_IN_DAYSis less than 14 days and a stream on it has not been consumed, Snowflake temporarily extends the retention period for that table, up to the value ofMAX_DATA_EXTENSION_TIME_IN_DAYS(default 14 days), regardless of edition. Setting that parameter to 0 disables the extension. SHOW STREAMSandDESCRIBE STREAMreport aSTALEflag and aSTALE_AFTERtimestamp: when the stream is predicted to become stale (or became stale, if in the past).
-- Snowflake SQL (not executed here)
SHOW STREAMS IN SCHEMA raw;
-- check the "stale" and "stale_after" columns
-- Allow up to 30 days before streams on this table go stale
ALTER TABLE raw.customers SET MAX_DATA_EXTENSION_TIME_IN_DAYS = 30;
Recovery: a stale stream must be recreated (CREATE OR REPLACE STREAM), which starts tracking from now. The changes between the old offset and now are not available through the stream, so you need a reconciliation: rebuild the target from the source, or compare source and target and apply the differences.
Pitfalls
- Suspending a task “for a few weeks” during a migration. Its stream goes stale silently.
- Lowering retention to save storage on a table that feeds streams. The extension is capped by
MAX_DATA_EXTENSION_TIME_IN_DAYS. - Monitoring task success but not stream staleness. A task with
WHEN SYSTEM$STREAM_HAS_DATAthat is never triggered does not fail; it just never runs.
In interviews
Define staleness, give the 14-day default extension, name STALE_AFTER, and describe recovery as recreate plus reconcile. Mention monitoring as the real prevention.
Task fundamentals
A task runs one SQL statement, a call to a stored procedure, or a block of Snowflake Scripting, either on a schedule, after another task, or when a stream has data.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE TASK refresh_daily_sales
WAREHOUSE = transform_wh
SCHEDULE = '60 MINUTE'
AS
INSERT OVERWRITE INTO analytics.daily_sales
SELECT order_date, SUM(amount) FROM analytics.orders GROUP BY order_date;
ALTER TASK refresh_daily_sales RESUME; -- tasks are created suspended
EXECUTE TASK refresh_daily_sales; -- run once now, for testing
Fundamentals:
- Compute: either a user-managed warehouse (
WAREHOUSE = ...) or Snowflake-managed serverless compute (omitWAREHOUSE). - Created suspended: a new or recreated task does nothing until
ALTER TASK ... RESUME. - Owner and privileges: tasks run with the privileges of the task’s owner role, not the user who created them. The owner needs the account-level
EXECUTE TASKprivilege (andEXECUTE MANAGED TASKfor serverless tasks), plusUSAGEon the warehouse and the privileges the SQL needs. - No overlap by default:
ALLOW_OVERLAPPING_EXECUTION = FALSEmeans a scheduled run is skipped if the previous run is still going. - Timeouts:
USER_TASK_TIMEOUT_MSlimits how long a run may take. - History:
INFORMATION_SCHEMA.TASK_HISTORY()(recent, near real time) andSNOWFLAKE.ACCOUNT_USAGE.TASK_HISTORY(a year, with latency).
Pitfalls
- Creating or replacing a task and forgetting to resume it.
CREATE OR REPLACEalso resets it to suspended. - An owner role that loses a privilege: the task fails at run time, not when created.
In interviews
Mention the two compute models, the “created suspended” behaviour and that tasks run as their owner role.
Scheduled tasks with CRON
SCHEDULE takes either an interval or a cron expression with a time zone:
| Form | Example | Meaning |
|---|---|---|
| Interval | SCHEDULE = '15 MINUTE' |
Every 15 minutes, counted from when the task is resumed |
| Cron | SCHEDULE = 'USING CRON 0 2 * * * UTC' |
02:00 UTC every day |
| Cron with time zone | SCHEDULE = 'USING CRON 30 6 * * MON-FRI Europe/London' |
06:30 London time on weekdays |
The cron fields are minute, hour, day of month, month and day of week, followed by an IANA time zone name.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE TASK nightly_rollup
WAREHOUSE = transform_wh
SCHEDULE = 'USING CRON 0 2 * * * UTC'
AS
CALL analytics.build_nightly_rollup();
Daylight saving time: with a local time zone, a schedule in the hour that is skipped or repeated when clocks change can run zero times or twice that day. Use UTC for anything that must run exactly once per day, or choose a time outside the change window.
Pitfalls
- Using an interval schedule for “every day at 2 a.m.”. Intervals drift from the resume time; use cron.
- Scheduling dozens of tasks at exactly midnight on one warehouse, then wondering why they queue.
In interviews
Know both forms, the time-zone suffix, and the daylight-saving caveat.
Task graphs (DAGs)
A task graph (formerly called a task tree) is a directed acyclic graph of tasks:
- the root task has the schedule (or trigger);
- child tasks declare predecessors with
AFTER, and run when all their predecessors have finished successfully in the same run; - an optional finalizer task runs after all other tasks in the graph finish, whether they succeeded or failed: use it for cleanup and alerts.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE TASK load_root
WAREHOUSE = transform_wh
SCHEDULE = 'USING CRON 0 3 * * * UTC'
AS
CALL staging.load_all();
CREATE OR REPLACE TASK build_customers
WAREHOUSE = transform_wh
AFTER load_root
AS
CALL analytics.build_dim_customer();
CREATE OR REPLACE TASK build_orders
WAREHOUSE = transform_wh
AFTER load_root
AS
CALL analytics.build_fct_orders();
CREATE OR REPLACE TASK build_marts
WAREHOUSE = transform_wh
AFTER build_customers, build_orders -- waits for both
AS
CALL analytics.build_marts();
CREATE OR REPLACE TASK cleanup_and_alert
WAREHOUSE = transform_wh
FINALIZE = load_root -- the finalizer for this graph
AS
CALL ops.cleanup_and_notify();
-- Resume every task in the graph, then the root
SELECT SYSTEM$TASK_DEPENDENTS_ENABLE('load_root');
Rules worth knowing:
- Only the root task has a schedule; a child with a schedule is rejected.
- A root can have only one finalizer; a finalizer cannot have children or a schedule.
- To change a graph, suspend the root first. Child tasks must be resumed for the graph to run them;
SYSTEM$TASK_DEPENDENTS_ENABLEresumes all dependents. - The documentation sets limits on the total number of tasks in a graph and on predecessors and children per task; check the current values if you build very large graphs.
- Tasks can pass small values to their children with
SYSTEM$SET_RETURN_VALUEandSYSTEM$GET_PREDECESSOR_RETURN_VALUE.
Pitfalls
- Modelling a complex, cross-system workflow (APIs, files, other platforms) as a task graph. Task graphs are good for SQL inside Snowflake; an orchestrator such as Airflow is better when steps live outside it.
- Forgetting the finalizer: without it, cleanup after a failed run never happens.
In interviews
Draw the graph, explain root, children with AFTER, and the finalizer, and say when you would use an external orchestrator instead.
Serverless tasks
Omit WAREHOUSE and the task runs on serverless compute that Snowflake sizes and manages:
-- Snowflake SQL (not executed here)
CREATE OR REPLACE TASK merge_customers_serverless
SCHEDULE = '5 MINUTE'
USER_TASK_MANAGED_INITIAL_WAREHOUSE_SIZE = 'XSMALL' -- first runs, before history exists
TARGET_COMPLETION_INTERVAL = '4 MINUTE' -- aim to finish within this time
AS
CALL analytics.merge_customers();
How it behaves:
- Snowflake uses the task’s run history to choose the compute size for later runs.
USER_TASK_MANAGED_INITIAL_WAREHOUSE_SIZEonly applies until enough history exists. TARGET_COMPLETION_INTERVALtells Snowflake how quickly runs should finish, so it can scale compute up for the run to meet it.SERVERLESS_TASK_MIN_STATEMENT_SIZEandSERVERLESS_TASK_MAX_STATEMENT_SIZEbound the sizes it may choose.- Billing is per second of compute actually used, at the serverless task rate listed in Snowflake’s Service Consumption Table, with no idle time and no warehouse minimum.
- The owner role needs
EXECUTE MANAGED TASK.
When to choose which:
| Prefer serverless | Prefer a warehouse |
|---|---|
| Short, frequent tasks (every few minutes) where warehouse resume minimums and idle time dominate | Long or heavy tasks that fully use a warehouse |
| Tasks that should finish on time without you sizing them | Tasks that can share an already-running warehouse with other work |
| Uneven workloads | When you need a resource monitor to cap spend (resource monitors cover warehouses, not serverless) |
Pitfalls
- Assuming serverless is always cheaper. The per-credit rate for serverless compute can differ from warehouse credits; compare actual costs in
SERVERLESS_TASK_HISTORYagainst warehouse metering.
In interviews
Explain how serverless tasks are sized (history plus initial size and target completion interval), how they are billed and when you would still use a warehouse.
The task plus stream pattern
The classic Snowflake CDC pipeline combines a stream and a task that runs only when the stream has changes.
Playbook: a stream and task CDC pipeline
Goal: keep analytics.dim_customer in sync with raw.customers, which is loaded continuously, including updates and deletes.
1. Create the stream on the source.
-- Snowflake SQL (not executed here)
CREATE OR REPLACE STREAM raw.customers_stream ON TABLE raw.customers;
2. Create the task that consumes it. Two options:
-- Snowflake SQL (not executed here)
-- Option A: scheduled, but skipped cheaply when there is nothing to do
CREATE OR REPLACE TASK analytics.apply_customer_changes
WAREHOUSE = transform_wh
SCHEDULE = '5 MINUTE'
WHEN SYSTEM$STREAM_HAS_DATA('raw.customers_stream')
AS
MERGE INTO analytics.dim_customer AS t
USING (
SELECT *
FROM raw.customers_stream
WHERE NOT (METADATA$ACTION = 'DELETE' AND METADATA$ISUPDATE)
QUALIFY ROW_NUMBER() OVER (
PARTITION BY customer_id, METADATA$ACTION ORDER BY updated_at DESC) = 1
) AS s
ON t.customer_id = s.customer_id
WHEN MATCHED AND s.METADATA$ACTION = 'DELETE' THEN DELETE
WHEN MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
UPDATE SET t.name = s.name, t.tier = s.tier, t.updated_at = s.updated_at
WHEN NOT MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
INSERT (customer_id, name, tier, updated_at)
VALUES (s.customer_id, s.name, s.tier, s.updated_at);
-- Option B: a triggered task (no SCHEDULE): runs when the stream has data
CREATE OR REPLACE TASK analytics.apply_customer_changes_triggered
WAREHOUSE = transform_wh
WHEN SYSTEM$STREAM_HAS_DATA('raw.customers_stream')
AS
CALL analytics.apply_customer_changes_proc();
The WHEN condition is evaluated by cloud services. If the stream is empty, the run is skipped and no warehouse resumes, so a 5-minute schedule costs nothing during quiet periods. A triggered task has no SCHEDULE; it runs when the stream has data, at most every 30 seconds by default (USER_TASK_MINIMUM_TRIGGER_INTERVAL_IN_SECONDS, minimum 10). For a serverless triggered task, TARGET_COMPLETION_INTERVAL is required.
3. Resume and test.
-- Snowflake SQL (not executed here)
ALTER TASK analytics.apply_customer_changes RESUME;
-- Make a change, then check the stream and the task history
UPDATE raw.customers SET tier = 'gold' WHERE customer_id = 2;
SELECT SYSTEM$STREAM_HAS_DATA('raw.customers_stream');
SELECT name, state, scheduled_time, completed_time, error_message
FROM TABLE(INFORMATION_SCHEMA.TASK_HISTORY(TASK_NAME => 'APPLY_CUSTOMER_CHANGES'))
ORDER BY scheduled_time DESC
LIMIT 10;
4. Monitor. Alert on task failures (next section) and on stream staleness (STALE_AFTER approaching), and reconcile row counts between source and target periodically.
Why it is exactly-once: the MERGE consumes the stream inside one statement. If it fails, nothing commits, the offset stays put, and the next run reprocesses the same changes.
Pitfalls
- Two tasks consuming one stream. Give each consumer its own stream.
- A
WHENcondition that references a stream the task does not consume: the stream is never emptied, so the task runs every time.
In interviews
Describe the four parts (stream, task with WHEN SYSTEM$STREAM_HAS_DATA, MERGE with the three action cases, monitoring) and the exactly-once argument. Compare it with dynamic tables, which can replace many such pipelines declaratively.
Task error handling
When a task run fails, the error is recorded in task history and the run is marked failed. The tools for handling failures:
| Tool | What it does |
|---|---|
SUSPEND_TASK_AFTER_NUM_FAILURES |
Automatically suspends a standalone task, or a graph’s root, after this many consecutive failed (or timed-out) runs. Default 10; 0 disables it. Set it on the root for a graph |
TASK_AUTO_RETRY_ATTEMPTS |
Set on the root: Snowflake retries a failed graph run from the failed task, up to this many times |
ERROR_INTEGRATION |
A notification integration (for example to Amazon SNS, Azure Event Grid or Google Pub/Sub) that receives error messages when a task fails |
| Finalizer task | Runs after the graph finishes, success or failure: send an email or post a message, clean up temporary tables |
TASK_HISTORY |
STATE, ERROR_CODE, ERROR_MESSAGE per run, for dashboards and alerts |
| Snowflake alerts | Scheduled conditions (for example “any failed task in the last hour”) that send notifications |
-- Snowflake SQL (not executed here)
ALTER TASK load_root SUSPEND;
ALTER TASK load_root SET
SUSPEND_TASK_AFTER_NUM_FAILURES = 3
TASK_AUTO_RETRY_ATTEMPTS = 2
ERROR_INTEGRATION = ops_error_notifications;
ALTER TASK load_root RESUME;
-- Failed runs in the last day
SELECT name, scheduled_time, error_code, error_message
FROM snowflake.account_usage.task_history
WHERE state = 'FAILED'
AND scheduled_time >= DATEADD('day', -1, CURRENT_TIMESTAMP())
ORDER BY scheduled_time DESC;
Inside the task, write steps so a retry is safe: consume streams in one statement, use MERGE or INSERT OVERWRITE rather than blind INSERT, and wrap multi-statement logic in a stored procedure with a transaction and an exception handler that logs context before re-raising.
Pitfalls
- A task that silently auto-suspends after repeated failures, followed by its stream going stale. Alert on suspension, not just on failure.
- Catching exceptions in a procedure and not re-raising, so the task reports success.
In interviews
Name the retry and auto-suspend parameters with the default of 10, describe notifications through an error integration or finalizer, and explain idempotent task design.
Practice questions
A stream on orders returned 1,000 rows. Your MERGE used WHERE METADATA$ACTION = 'INSERT' and committed. How many rows will the stream return now?
Zero, until new changes arrive. Consuming a stream in a committed DML statement advances the offset past all changes as of the start of the transaction, regardless of any WHERE filter. The DELETE rows you filtered out are no longer available through that stream.
A row is inserted, updated twice and then deleted between two task runs. What does a standard stream show? An append-only stream?
The standard stream shows nothing for that row, because it returns net changes and the row did not exist at either end of the interval. The append-only stream shows one INSERT row with the values as originally inserted, because it records inserts and ignores updates and deletes.
Your task was suspended for three weeks during a migration. When you resume it, the MERGE fails because the stream is stale. What happened and how do you recover?
The stream’s offset fell outside the source table’s retention period. Snowflake extends retention for unconsumed streams only up to MAX_DATA_EXTENSION_TIME_IN_DAYS (default 14 days), so after three weeks the history was purged. Recreate the stream with CREATE OR REPLACE STREAM (it starts tracking from now) and reconcile the target: rebuild it from the source, or compare and apply differences. To prevent it, monitor STALE_AFTER and raise MAX_DATA_EXTENSION_TIME_IN_DAYS on that table before planned pauses.
Why does a task scheduled every minute with WHEN SYSTEM$STREAM_HAS_DATA(...) not cost a warehouse minute every minute?
The WHEN condition is evaluated in the cloud services layer before the task needs compute. If the stream has no data, the run is skipped and the warehouse is not resumed. Only runs with data consume warehouse credits (with the usual 60-second minimum on resume).
Two teams want to consume changes from the same table on different schedules. What do you do?
Create one stream per consumer. A stream has a single offset; if both teams used the same stream, whichever consumed first would advance the offset and the other would miss those changes.
When would you choose a serverless task over a warehouse task?
For short, frequent or uneven tasks where a warehouse would spend much of its billed time idle or paying the 60-second resume minimum, and when you want Snowflake to size compute to meet a completion target. Use a warehouse for long, heavy work that fully uses it, for tasks that can share an already-running warehouse, or when you need resource monitors to cap spend. Compare real costs in SERVERLESS_TASK_HISTORY and warehouse metering.
How do you alert someone when any task in a graph fails?
Options: set ERROR_INTEGRATION on the root to a notification integration so failures are pushed to a cloud messaging service; add a finalizer task that checks the graph’s run results and sends a notification; or create a Snowflake alert on TASK_HISTORY for failed states. Also alert when the root is auto-suspended after SUSPEND_TASK_AFTER_NUM_FAILURES consecutive failures.
Key takeaways
- A stream stores an offset into a table’s history and returns changed rows with
METADATA$ACTION,METADATA$ISUPDATEandMETADATA$ROW_ID; it does not copy data. - The offset moves only when the stream is consumed in a committed DML statement, all changes at once; rollbacks leave it in place, which gives exactly-once processing.
- Standard streams return net inserts, updates and deletes; append-only streams return inserts only; insert-only streams are for external and externally managed Iceberg tables.
- Streams go stale when their offset leaves the retention period, extended for unconsumed streams up to
MAX_DATA_EXTENSION_TIME_IN_DAYS(default 14); monitorSTALE_AFTER. - Tasks run on a warehouse or serverless compute, on an interval, a cron schedule, after other tasks, or when a stream has data; they are created suspended.
- Combine
WHEN SYSTEM$STREAM_HAS_DATAwith aMERGEthat handles the three action cases, and protect it with retries, auto-suspension, notifications and a finalizer.
Progress is saved in this browser only. No account needed.