CertSafari
    Snowflake SnowPro Core Certification (COF-C03)· Lessons

    Domain 3 · Lesson 11/19

    Streams, Tasks and Dynamic Tables: Building CDC Pipelines in Snowflake

    Perform automated data ingestion

    16 min read
    6% of exam
    6 sources
    Published 5 Oct 2026
    Docs as of 4 Oct 2026

    What you will be able to do

    • Explain how a stream's offset works and what does and doesn't advance it
    • Pick the right stream type: standard, append-only or insert-only
    • Configure a task's compute model, schedule (interval or CRON) and stream trigger
    • Chain tasks into a task graph with CREATE TASK .. AFTER
    • Define a dynamic table with a target lag and refresh mode as the declarative alternative to streams and tasks

    1.Streams: a bookmark on a table's changes

    Once data has landed in a table, a pipeline needs to know what has changed since it last ran. A stream records DML changes made to a source object: inserts (including COPY INTO), updates and deletes. This is called change data capture (CDC). You can create streams on standard tables, views, directory tables, dynamic tables, Iceberg tables, event tables and external tables.

    A stream holds no table data of its own. It stores an offset, which is a point in the source's version history, and builds the change records from the source's versioning history plus hidden change-tracking columns. Those columns are added to the table when its first stream is created. Snowflake compares a stream to a bookmark. Querying it returns every change committed after the offset, up to the current time. The rows have the same columns as the source, plus METADATA$ACTION (INSERT or DELETE), METADATA$ISUPDATE and METADATA$ROW_ID. An UPDATE shows up as a DELETE/INSERT pair with METADATA$ISUPDATE = TRUE.

    The offset rule is the part most often tested. Selecting from a stream does not move the offset. Many queries can read the same changes independently. The offset only advances when the stream is used in a DML transaction, such as INSERT … SELECT, CTAS or COPY INTO location, and that transaction commits. If the transaction does not commit, the offset stays where it was. Inside an explicit BEGIN … COMMIT transaction, the stream is locked and every statement sees the same change set (repeatable read).

    Checkpoint 1 of 9· Check yourself

    Which of these moves a stream's offset forward?

    Streams come in different types, depending on what changes they record. A standard stream joins inserted and deleted rows to work out the net change, so a row that was inserted and then deleted between two offsets doesn't appear at all. An append-only stream records inserts only.

    Standard vs append-only streams
    Stream typeWhat it tracksSupported on
    Standard (delta)All DML: inserts, updates, deletes, including table truncatesStandard tables, dynamic tables, Iceberg tables, directory tables, views
    Append-onlyRow inserts only; updates, deletes and truncates are not capturedStandard tables, dynamic tables, Snowflake-managed Iceberg tables, views

    External tables use a third type, the insert-only stream. It tracks records added to the external table's metadata, so new files show up as new rows. You create one by adding INSERT_ONLY = TRUE:

    An insert-only stream on an external tablesql
    CREATE STREAM my_ext_table_stream ON EXTERNAL TABLE my_ext_table INSERT_ONLY = TRUE;

    Staleness. A stream becomes stale when its offset falls outside the data retention period of its source table. A stale stream can no longer return its unconsumed change records, and you must recreate it with CREATE STREAM. To help, if a table's retention is under 14 days and a stream hasn't been consumed, Snowflake temporarily extends the retention to the stream's offset, up to 14 days by default. The MAX_DATA_EXTENSION_TIME_IN_DAYS parameter sets that maximum. Once the stream is consumed, retention reverts to the table's default.

    SHOW STREAMS and DESCRIBE STREAM show a STALE_AFTER timestamp, which is the last consumption time plus the larger of DATA_RETENTION_TIME_IN_DAYS and MAX_DATA_EXTENSION_TIME_IN_DAYS. To prevent staleness, consume the stream in a DML statement before STALE_AFTER. Calling SYSTEM$STREAM_HAS_DATA on a stream also prevents staleness, provided the stream is empty and the function returns FALSE. Recreating a table with CREATE OR REPLACE TABLE drops its history and makes any stream on it stale.

    Advancing without consuming. To move the offset to the current table version without processing the changes, either recreate the stream with CREATE OR REPLACE STREAM, or insert the stream's data into a temporary table with a WHERE clause that filters out every row (for example WHERE 0 = 1).

    Checkpoint 2 of 9· Check yourself

    A table has 1-day data retention and a stream on it has not been consumed for 5 days. By default, what stops the stream from going stale?

    Checkpoint 3 of 9· Exam question

    A retail company loads thousands of small CSV files (average 200 KB) into Snowflake every few seconds using Snowpipe with auto-ingest configured against an Amazon S3 event notification through SQS. The data engineering team notices ingestion costs are unexpectedly high and file processing latency is inconsistent. What is the most effective way to reduce Snowpipe costs while maintaining near-real-time ingestion?

    Sources12

    2.Tasks: compute, schedules and stream triggers

    A stream records what changed, and a task runs the work that processes it. A task can run SQL or a stored procedure written in JavaScript, Python, Java, Scala or Snowflake Scripting. It runs either on a schedule or when an event occurs.

    Compute. There are two models. In a *serverless* task you leave out the WAREHOUSE parameter, and Snowflake sizes compute from recent runs. You can bound that sizing with SERVERLESS_TASK_MIN_STATEMENT_SIZE (default XSMALL) and SERVERLESS_TASK_MAX_STATEMENT_SIZE (default XXLARGE), and XXLARGE is the upper limit. In a *user-managed* task you name a warehouse yourself. For both models, the documentation says the role that runs the task must have the global EXECUTE MANAGED TASK privilege.

    Serverless task with bounded compute size, running every 30 secondssql
    CREATE TASK SCHEDULED_T2
      SCHEDULE='30 SECONDS'
      SERVERLESS_TASK_MIN_STATEMENT_SIZE='SMALL'
      SERVERLESS_TASK_MAX_STATEMENT_SIZE='LARGE'
      AS SELECT 1;
    Serverless vs user-managed tasks
    FactorServerless tasksUser-managed tasks
    Workload fitUnder-utilized warehouses; tasks with relatively stable runsFully utilized warehouses with multiple concurrent tasks; unpredictable loads
    Schedule adherenceRecommended when adherence is highly important; compute grows if a run exceeds the intervalRecommended when adherence is less important
    BillingActual compute resource usageWarehouse size, 60-second minimum each time the warehouse resumes
    Size ceilingEquivalent to XXLARGEAny warehouse size you choose

    Schedules. Set SCHEDULE to an interval such as '60 MINUTES', or to 'USING CRON …' with a time zone to run at a specific time or day. Only one instance of a scheduled task runs at a time. If the previous run is still going when the next one is due, that scheduled run is skipped.

    A CRON schedule: every Sunday at 3:07 a.m. Pacificsql
    CREATE TASK task_sunday_3_07_am_pacific_time_zone
      SCHEDULE='USING CRON 7 3 * * SUN America/Los_Angeles'  -- Use a random minute such as 7
    AS SELECT 1;

    Triggered tasks. A WHEN clause with SYSTEM$STREAM_HAS_DATA makes the task run only when the stream has changes. This avoids polling a source whose data arrives unpredictably. You can also combine SCHEDULE with WHEN, so the task checks the stream on a schedule and runs only when there is data. Serverless triggered tasks must have a TARGET_COMPLETION_INTERVAL.

    Lifecycle. A new task is created suspended. EXECUTE TASK runs it once, for testing. ALTER TASK … RESUME lets it follow its schedule or respond to triggers.

    Checkpoint 4 of 9· Fill the gap

    Which function makes this task run only when the stream has new change records?

    CREATE TASK triggered_task_stream
      WHEN  ? ('orders_stream')
      AS
        INSERT INTO completed_promotions
        SELECT order_id, order_total, order_time, promotion_id
        FROM orders_stream;

    Checkpoint 5 of 9· Exam question

    An engineer creates a Snowpipe object with `AUTO_INGEST = TRUE` referencing an internal stage that points at an Amazon S3 bucket, but new files landing in the bucket are never automatically loaded into the target table. The pipe was created successfully and the stage lists the files when queried manually. What is the most likely cause of the missing automatic ingestion?

    Sources3

    3.Task graphs: chaining steps

    Real pipelines usually have several steps. A task graph, also called a DAG, is made of a root task and the child tasks that depend on it. Dependencies run from start to finish with no loops. You create the root with CREATE TASK and attach each child with CREATE TASK .. AFTER, naming its parent or parents. Only the root task's schedule controls when the graph runs, and child tasks run in the order the graph defines.

    The graph shape controls concurrency. Children of the same parent run in parallel. A task with several parents waits until all of them have completed successfully. To make steps run one after another, make each step the child of the step before it. You can also add a finalizer, an optional last task that does cleanup after every other task has finished.

    Checkpoint 6 of 9· Check yourself

    Root task R loads raw data. Transformations T1 and T2 are both created with AFTER R. How do T1 and T2 run once R completes?

    Limits. A task graph can hold at most 1000 tasks, and a single task can have at most 100 parent tasks and 100 child tasks.

    Finalizer syntax. You create a finalizer with CREATE TASK … FINALIZE = <root task>. It runs after all other tasks in the graph complete or fail to complete, so it suits cleanup and notifications. Each root task can have only one finalizer, and a finalizer can be tied to only one root. A finalizer cannot have child tasks and cannot have a schedule.

    A finalizer task attached to the root tasksql
    CREATE TASK task_finalizer
      FINALIZE = task_root
      AS SELECT 1;

    Starting a graph. Either resume each child task you want in the run (including the finalizer) and then the root task with ALTER TASK … RESUME, or resume every task at once by calling SYSTEM$TASK_DEPENDENTS_ENABLE with the root task's name. To run the graph once for testing, resume the child tasks and run EXECUTE TASK on the root.

    Checkpoint 7 of 9· Check yourself

    Which statement about a finalizer task is correct?

    Checkpoint 8 of 9· Exam question

    A data engineer needs to reload a corrected version of a file named `sales_2026_09_10.csv` that Snowpipe already loaded successfully into the SALES table. The engineer overwrites the file in the stage with the corrected data using the same file name and waits for Snowpipe to reprocess it automatically. Why does the corrected data never reach the table?

    Sources4

    4.Dynamic tables: the declarative alternative

    With streams and tasks, you write the change logic and the scheduling yourself. A dynamic table only needs a SELECT query and a target freshness. Snowflake materializes the query result and keeps it up to date. The Snowpipe Streaming documentation points readers who want SQL-native streaming to dynamic tables and to streams with tasks.

    A dynamic table with a 10-minute target lag and incremental refreshsql
    CREATE OR ALTER DYNAMIC TABLE dt_orders
        TARGET_LAG = '10 minutes'
        WAREHOUSE = transform_wh
        REFRESH_MODE = INCREMENTAL
    AS
        SELECT
            order_id,
            customer_id,
            order_date,
            TRIM(UPPER(product_name)) AS product_name,
            quantity,
            unit_price,
            quantity * unit_price AS line_total,
            order_status
        FROM raw_orders
        WHERE order_status != 'returned';

    TARGET_LAG sets how far the table may fall behind its base tables. It is a target, not a guarantee: if refreshes take longer than expected, the actual lag can be larger. On intermediate tables, TARGET_LAG = DOWNSTREAM means they refresh only when a downstream table needs fresh data. REFRESH_MODE can be INCREMENTAL (only rows that changed), FULL (the whole result), AUTO (Snowflake picks at creation), ADAPTIVE (incremental, but reinitializes after large upstream changes) or CUSTOM_INCREMENTAL (your own DML).

    You don't declare dependencies between dynamic tables. When one dynamic table reads from another, Snowflake works out the dependency graph from the queries and refreshes the tables in dependency order, using a consistent snapshot. SCHEDULER = DISABLE switches off automatic refreshes so you can trigger them by hand or from a tool such as dbt or Airflow.

    That is the practical contrast with streams and tasks. With a dynamic table you declare the result and the lag, and Snowflake handles dependencies and refresh scheduling. With streams and tasks you script the change processing, the schedule or trigger, and the order of the steps yourself.

    Checkpoint 9 of 9· Check yourself

    What does TARGET_LAG = DOWNSTREAM do on an intermediate dynamic table?

    Sources56

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 1.Querying a stream, at least inside an explicit transaction, consumes its changes and moves the offset forward.Why is that wrong?

      A SELECT never moves the offset forward. The stream has to be consumed by a DML statement in a transaction that commits.

      Covered in Streams: a bookmark on a table's changes

    2. 2.A task starts following its SCHEDULE as soon as CREATE TASK succeeds.Why is that wrong?

      New tasks are created suspended. ALTER TASK … RESUME starts the schedule, and EXECUTE TASK runs the task once.

      Covered in Tasks: compute, schedules and stream triggers

    3. 3.A serverless task can scale to any warehouse size the workload needs.Why is that wrong?

      Serverless tasks top out at the equivalent of XXLARGE. Workloads that need more must use a user-managed task with a larger warehouse.

      Covered in Tasks: compute, schedules and stream triggers

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 1.
      “a stream itself does not contain any table data”
      ↩︎ Streams: a bookmark on a table's changes
      “Updates to rows in the source object are represented as a pair of DELETE and INSERT records in the stream”
      ↩︎ Streams: a bookmark on a table's changes
      “An append-only stream exclusively tracks row inserts.”
      ↩︎ Streams: a bookmark on a table's changes
      “A stream becomes stale when its offset falls outside of the data retention period for its source table”
      ↩︎ Streams: a bookmark on a table's changes
      “Recreate the stream (using the CREATE OR REPLACE STREAM syntax).”
      ↩︎ Streams: a bookmark on a table's changes
      “Querying a stream alone does not advance its offset, even within an explicit transaction; the stream contents must be consumed in a DML statement.”
      ↩︎ Exam trap 1
      “A stream advances the offset only when it is used in a DML transaction.”
      ↩︎ Checkpoint
      “The retention period is extended to the stream’s offset, up to a maximum of 14 days by default, regardless of your Snowflake edition.”
      ↩︎ Checkpoint
    2. 2.
      “track the records added to the external table metadata”
      ↩︎ Streams: a bookmark on a table's changes
    3. 3.
      “Serverless tasks: Snowflake predicts resources that are needed and assigns them automatically.”
      ↩︎ Tasks: compute, schedules and stream triggers
      “If a task is still running when the next scheduled run time occurs, then that scheduled time is skipped.”
      ↩︎ Tasks: compute, schedules and stream triggers
      “A target completion interval is required for serverless triggered tasks.”
      ↩︎ Tasks: compute, schedules and stream triggers
      “The role that runs the task must have the global EXECUTE MANAGED TASK privilege.”
      ↩︎ Tasks: compute, schedules and stream triggers
      “When a task is created, it starts as suspended.”
      ↩︎ Exam trap 2
      “The maximum compute size for a serverless task is equivalent to an XXLARGE virtual warehouse.”
      ↩︎ Exam trap 3
    4. 4.
      “Create a root task using CREATE TASK, then create child tasks using CREATE TASK .. AFTER to select the parent tasks.”
      ↩︎ Task graphs: chaining steps
      “When a task has multiple parents, the task waits for all preceding tasks to successfully complete before starting.”
      ↩︎ Task graphs: chaining steps
      “A task graph is limited to a maximum of 1000 tasks.”
      ↩︎ Task graphs: chaining steps
      “A single task can have a maximum of 100 parent tasks and 100 child tasks.”
      ↩︎ Task graphs: chaining steps
      “Resume all of the tasks in a task graph at once by calling SYSTEM$TASK_DEPENDENTS_ENABLE”
      ↩︎ Task graphs: chaining steps
      “When multiple child tasks have the same parent, the child tasks run in parallel.”
      ↩︎ Checkpoint
      “Each root task can have only one finalizer task, and a finalizer task can be associated with only one root task.”
      ↩︎ Checkpoint
    5. 5.
      “A dynamic table materializes the results of a SELECT query and keeps them up to date.”
      ↩︎ Dynamic tables: the declarative alternative
      “actual lag can exceed the target when refreshes take longer than expected”
      ↩︎ Dynamic tables: the declarative alternative
      “Snowflake infers the dependency graph automatically from the queries you write.”
      ↩︎ Dynamic tables: the declarative alternative
      “so they refresh only when their downstream dependents need fresh data”
      ↩︎ Checkpoint
    6. 6.

    Ready to test yourself?

    Practise the 24 questions on this subdomain.

    Spotted a mistake, or was something unclear? Tell us.