CertSafari
    Snowflake SnowPro Advanced: Data Engineer (DEA-C02)· Lessons

    Domain 1 · Lesson 4/22

    Streams, tasks, dynamic tables and programmatic pipelines

    Design, build, and troubleshoot continuous data pipelines.

    18 min read
    4% of exam
    13 sources
    Published 5 Oct 2026
    Docs as of 4 Oct 2026

    What you will be able to do

    • Explain how a stream's offset advances and what that means for consuming change data
    • Build scheduled and stream-triggered tasks, and choose between serverless and user-managed compute
    • Choose between dynamic tables, materialized views and streams with tasks
    • Use UDFs, Snowflake Scripting blocks and the SQL API to automate pipeline steps

    1.Streams: change data capture with an offset

    After data lands, a continuous pipeline has to process only what has changed. A stream records DML changes (inserts, including COPY INTO, plus updates and deletes) on a source object. You can then query those changes as a change table. Streams can be created on standard and shared tables, views (including secure views), directory tables, dynamic tables, Iceberg tables (with limitations), event tables and external tables.

    A stream doesn't contain any table data. It stores only an offset, which works like a bookmark between two table versions. It then returns change records from the source object's version history. When you create the first stream on a table, Snowflake adds hidden change-tracking columns to that table. Views are an exception: for a stream on a view, you must enable change tracking explicitly on the view and on its underlying tables.

    The offset advances only when a DML transaction consumes the stream. CTAS and COPY INTO location count as DML here. Running a plain SELECT on a stream, even inside an explicit transaction, leaves the offset where it is. If several statements must see the same change records, wrap them in BEGIN … COMMIT. This locks the stream, and because streams use repeatable read isolation, every statement in the transaction sees the same set of changes.

    Checkpoint 1 of 9· Exam question

    A retail company lands clickstream JSON files in an S3 bucket every few seconds and wants them loaded into a raw table within about a minute of arrival, with no external service polling S3 or calling Snowflake. Which design meets this requirement?

    Checkpoint 2 of 9· Check yourself

    An analyst runs SELECT * FROM orders_stream three times to inspect changes, and then a task runs INSERT … SELECT FROM orders_stream. When does the offset advance?

    Sources1

    2.Tasks: schedules and stream triggers

    A task runs a SQL command or a stored procedure on a fixed schedule or when an event happens. Stored procedures can be written in JavaScript, Python, Java, Scala or Snowflake Scripting. A schedule can be an interval such as '10 SECONDS', or a CRON expression with a time zone. Only one instance of a scheduled task runs at a time. If a run is still going when the next one is due, that next run is skipped.

    A triggered task uses WHEN SYSTEM$STREAM_HAS_DATA(...) and runs when the stream has new data. When data arrives unpredictably, this avoids frequent polling and processes new data sooner. You can also add a SCHEDULE to a triggered task, so the stream check runs at that interval. Groups of tasks can be chained into task graphs that run steps in parallel or in sequence.

    A triggered task that consumes a stream whenever it has datasql
    CREATE TASK triggered_task_stream
      WHEN SYSTEM$STREAM_HAS_DATA('orders_stream')
      AS
        INSERT INTO completed_promotions
        SELECT order_id, order_total, order_time, promotion_id
        FROM orders_stream;

    A new task starts in the suspended state. Use EXECUTE TASK to test it once, and ALTER TASK … RESUME to start its schedule or event detection.

    For multi-step pipelines, build a task graph with parent-child dependencies, fan-out and fan-in patterns, and finalizer tasks. When you create a task, you also define what happens when it fails. For example, SUSPEND_TASK_AFTER_NUM_FAILURES = 3 suspends a task after three failures. To diagnose problems such as an auto-suspended task, or a task that works interactively but fails under the scheduler, inspect TASK_HISTORY and the owner role's privileges.

    Checkpoint 3 of 9· Exam question

    An on-premises ETL tool writes export files to a Snowflake internal stage rather than cloud storage, so no cloud storage event notification service is available to trigger loading. The team still wants each new file loaded within seconds of the ETL job finishing its write. Which mechanism should they use?

    Checkpoint 4 of 9· Fill the gap

    Which function makes this hourly task act only when the stream has new data?

    CREATE TASK triggered_task_stream
      SCHEDULE = '1 HOUR'
      WHEN  ? ('orders_stream')
      AS SELECT 1;

    Sources23

    3.Serverless or user-managed task compute

    Every task needs compute. If you leave out the WAREHOUSE parameter, the task is serverless. Snowflake predicts the compute it needs from recent runs, within SERVERLESS_TASK_MIN_STATEMENT_SIZE and SERVERLESS_TASK_MAX_STATEMENT_SIZE. The defaults are XSMALL and XXLARGE. If you include WAREHOUSE, the task is user-managed, and you control the compute size directly. In both cases, the role that runs the task needs the global EXECUTE MANAGED TASK privilege.

    TARGET_COMPLETION_INTERVAL tells Snowflake to scale a serverless task so it finishes sooner. It is required for serverless triggered tasks. When a task has already reached its maximum size, Snowflake ignores the interval.

    A serverless task bounded between SMALL and LARGE computesql
    CREATE TASK SCHEDULED_T2
      SCHEDULE='30 SECONDS'
      SERVERLESS_TASK_MIN_STATEMENT_SIZE='SMALL'
      SERVERLESS_TASK_MAX_STATEMENT_SIZE='LARGE'
      AS SELECT 1;
    Serverless compared to user-managed tasks
    FactorServerlessUser-managed
    DefinitionOmit WAREHOUSEInclude WAREHOUSE
    Best fitUnder-utilized warehouses, stable runs, strict schedule adherenceFully utilized warehouses, unpredictable loads
    BillingActual compute resource usageWarehouse size, 60-second minimum per resume
    Size ceilingEquivalent to XXLARGEAny warehouse size you choose

    Checkpoint 5 of 9· Exam question

    A data engineer suspects a Snowpipe pipe silently stopped loading files two days ago after a schema drift caused repeated COPY errors, and finance is now asking why Snowpipe credit consumption dropped. Where should the engineer look first to confirm both the load failures and the associated compute usage?

    Checkpoint 6 of 9· Check yourself

    A transformation task needs more compute than an XXLARGE warehouse. How should it be defined?

    Sources2

    4.Dynamic tables and materialized views

    Streams and tasks require you to write the orchestration yourself. A dynamic table is the declarative alternative. You write a SELECT and set a TARGET_LAG, and Snowflake works out the dependency graph and keeps the results fresh. Each refresh is applied atomically. When dynamic tables read from one another, Snowflake refreshes them in dependency order against a consistent snapshot.

    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 is a goal, not a guarantee: the actual lag can exceed it if a refresh runs long. Setting TARGET_LAG = DOWNSTREAM on an intermediate table makes it refresh only when the tables that depend on it need fresh data. The refresh modes are INCREMENTAL, FULL, AUTO (chosen at creation), ADAPTIVE and CUSTOM_INCREMENTAL. Costs come from warehouse compute, cloud services and storage.

    Dynamic tables don't fit workloads that need data fresher than 60 seconds, the minimum target lag, or that need strictly guaranteed refresh timing. They also can't use stored procedures or external functions in their definition. For those cases, use streams with tasks.

    A materialized view solves a different problem. It precomputes a query result to speed up frequent or expensive aggregation, projection and selection on large data sets. In other words, it makes reads faster, whereas a dynamic table is a step in a transformation pipeline.

    Materialized views have trade-offs and limits worth knowing:

    - Always current. A background service maintains the view automatically. If a query runs before the view is up to date, Snowflake updates it or reads the newer data from the base table. A dynamic table has a minimum lag of 1 minute, so use a materialized view when you need always-current data with no lag. - Cost. Keeping the view up to date uses compute, and the stored results use storage. A materialized view is worth it only when the saved work outweighs those costs. - Single base table. Materialized views accelerate repeated queries against a single base table, with no joins. Dynamic tables are the choice for joins, aggregations and multi-step pipelines. - Transparent use. The optimizer rewrites queries to use a materialized view automatically. That doesn't happen with dynamic tables. - When to create one. Create one when the results don't change often, are used often, and the query consumes a lot of resources. Otherwise a regular view is the better fit.

    In short, use a materialized view for read performance on one table, and a dynamic table to build a pipeline.

    Checkpoint 7 of 9· Check yourself

    A pipeline must reflect source changes within 20 seconds and calls a stored procedure during transformation. Which approach do the sources rule out?

    Sources4567

    5.UDFs, Snowflake Scripting and Notebooks

    Pipeline logic that you reuse can live in a UDF. CREATE FUNCTION defines a function that returns either scalar or tabular results. Its handler is written in a supported language. Depending on the language, the handler code is either written inline in the statement or kept on a stage. CREATE OR ALTER FUNCTION creates the function if it doesn't exist, or alters it if it does.

    UDF handler languages and where the handler code can live
    LanguageHandler location
    JavaIn-line or staged
    JavaScriptIn-line
    PythonIn-line or staged
    ScalaIn-line or staged
    SQLIn-line

    For procedural steps, such as branching, variables and error handling, use Snowflake Scripting. You write the code in a block, which tasks can run inside stored procedures. Only BEGIN and END are required. DECLARE and EXCEPTION are optional. A block's BEGIN isn't the same as the BEGIN that starts a transaction, so Snowflake recommends starting transactions with BEGIN TRANSACTION. This matters when a block also has to consume a stream inside a transaction.

    Structure of a Snowflake Scripting blocksql
    DECLARE
      -- (variable declarations, cursor declarations, etc.) ...
    BEGIN
      -- (Snowflake Scripting and SQL statements) ...
    EXCEPTION
      -- (statements for handling exceptions) ...
    END;

    The exam guide also lists using Notebooks to run pipelines of stored procedures for ingestion. None of the sources for this lesson covers Notebooks, so this lesson can't describe how they work. Check the Notebooks documentation directly for that bullet.

    Checkpoint 8 of 9· Check yourself

    A developer writes a Snowflake Scripting block and wants an explicit transaction around two DML statements that consume a stream. How should the transaction start?

    Sources892

    6.Driving pipelines through the SQL API

    External orchestrators can run pipeline SQL through the Snowflake SQL API, a REST API for running queries and most DDL and DML. It has three endpoints under /api/v2/statements/: one submits statements, one checks a statement's status by its statement handle, and one cancels a statement. You submit with a POST that sends the statement and, optionally, a timeout, bind variables, and the warehouse, database, schema and role. If you leave out the context fields, the API uses the user's defaults. These values are case-sensitive and must match what SHOW returns.

    The API authenticates with OAuth or a key pair. In the documented curl example, the request carries a JWT that you generated in an Authorization: Bearer header.

    Request body for POST /api/v2/statements with a bind variablejson
    {
      "statement": "select * from T where c1=?",
      "timeout": 60,
      "database": "TESTDB",
      "schema": "TESTSCHEMA",
      "warehouse": "TESTWH",
      "role": "TESTROLE",
      "bindings": {
        "1": {
          "type": "FIXED",
          "value": "123"
        }
      }
    }

    Your client must handle asynchronous responses. If async=true isn't set, the API returns results only when the statement finishes within 45 seconds. Otherwise it returns a statement handle and HTTP 202, and you poll with GET. A 408 means the statement exceeded its timeout and was cancelled. A 422 means it failed.

    The API has limits that matter for pipelines. PUT and GET aren't supported, so you can't use it to upload files to a stage. BEGIN, COMMIT, ROLLBACK, USE, ALTER SESSION and temporary-object creation work only inside a multi-statement request. AUTOCOMMIT must be TRUE at the statement level. Some Python and Java/Scala stored procedures that return Arrow result sets can fail.

    Checkpoint 9 of 9· Check yourself

    A client posts a long-running MERGE to /api/v2/statements without the async parameter, and it takes three minutes. What does the client receive first?

    Sources10111213

    Exam traps

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

    1. 1.Selecting from a stream marks those changes as consumed.Why is that wrong?

      Only a committed DML transaction that uses the stream advances its offset. Plain queries can read the same changes as many times as needed.

      Covered in Streams: change data capture with an offset

    2. 2.A serverless triggered task needs only the WHEN condition; Snowflake picks the timing.Why is that wrong?

      Serverless triggered tasks must set TARGET_COMPLETION_INTERVAL.

      Covered in Serverless or user-managed task compute

    3. 3.Setting TARGET_LAG guarantees the dynamic table is never staler than that value.Why is that wrong?

      Target lag is a goal. Actual lag can be longer when refreshes run long, and the minimum target is 60 seconds.

      Covered in Dynamic tables and materialized views

    Practise it for real

    Turn a polling job into a stream-triggered task and bring it online safely.

    1. 1.Create a user-managed triggered task: CREATE TASK triggered_task_stream WAREHOUSE = transform_wh WHEN SYSTEM$STREAM_HAS_DATA('orders_stream') AS INSERT INTO completed_promotions SELECT ... FROM orders_stream; (transform_wh is a warehouse you own.)

      Why: With a WHEN condition, the task runs only when the stream has new data, so it doesn't poll for nothing. Including WAREHOUSE makes it user-managed, so no TARGET_COMPLETION_INTERVAL is needed.

      You should see: The task exists in a suspended state, because every new task starts suspended.

    2. 2.Run EXECUTE TASK triggered_task_stream;

      Why: Manually testing one run is the documented step before you resume a task.

      You should see: A single run that inserts any pending stream rows into completed_promotions and advances the stream offset.

    3. 3.Run ALTER TASK triggered_task_stream RESUME;

      Why: Resuming lets the task detect stream events continuously.

      You should see: The task now runs whenever orders_stream has new data.

    4. 4.Review the task's costs and history, then refine it with ALTER TASK.

      Why: Monitoring costs and then refining with ALTER TASK are the final steps in the documented task workflow.

      You should see: Settings adjusted based on observed runs.

    Stuck? Get a nudge

    If you leave out WAREHOUSE to make this task serverless, add a TARGET_COMPLETION_INTERVAL. Serverless triggered tasks require one.

    Sources

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

    1. 1.
      “for streams on views, change tracking must be enabled explicitly for the view and underlying tables”
      ↩︎ Streams: change data capture with an offset
      “To ensure multiple statements access the same change records in the stream, surround them with an explicit transaction statement (BEGIN .. COMMIT).”
      ↩︎ Streams: change data capture with an offset
      “A stream advances the offset only when it is used in a DML transaction.”
      ↩︎ Exam trap 1
      “Querying a stream alone does not advance its offset, even within an explicit transaction”
      ↩︎ Checkpoint
    2. 2.
      “it eliminates frequent polling of the source when new data arrival is unpredictable”
      ↩︎ Tasks: schedules and stream triggers
      “If a task is still running when the next scheduled run time occurs, then that scheduled time is skipped.”
      ↩︎ Tasks: schedules and stream triggers
      “When a task is created, it starts as suspended.”
      ↩︎ Tasks: schedules and stream triggers
      “For complex workflows, you can create sequences of tasks called task graphs.”
      ↩︎ Tasks: schedules and stream triggers
      “A target completion interval is required for serverless triggered tasks.”
      ↩︎ Serverless or user-managed task compute
      “The maximum compute size for a serverless task is equivalent to an XXLARGE virtual warehouse.”
      ↩︎ Serverless or user-managed task compute
      “Tasks can run SQL commands and stored procedures that use supported languages and tools, including JavaScript, Python, Java, Scala, and Snowflake scripting.”
      ↩︎ UDFs, Snowflake Scripting and Notebooks
      “A target completion interval is required for serverless triggered tasks.”
      ↩︎ Exam trap 2
      “If a task workload requires a larger warehouse, create a user-managed task with a warehouse of the required size.”
      ↩︎ Checkpoint
    3. 3.
      “Build a multi-step pipeline as a task graph with parent-child dependencies, fan-out and fan-in patterns, and finalizer tasks”
      ↩︎ Tasks: schedules and stream triggers
    4. 4.
      “A dynamic table materializes the results of a SELECT query and keeps them up to date.”
      ↩︎ Dynamic tables and materialized views
      “Need stored procedures or external functions in the definition.”
      ↩︎ Dynamic tables and materialized views
      “actual lag can exceed the target when refreshes take longer than expected”
      ↩︎ Exam trap 3
      “Require data fresher than 60 seconds (the minimum target lag) or strictly guaranteed refresh timing.”
      ↩︎ Checkpoint
    5. 5.
      “Because the result is pre-computed, querying a materialized view is faster than executing a query against the base table of the view.”
      ↩︎ Dynamic tables and materialized views
      “Incurs compute to keep up to date. Consumes storage.”
      ↩︎ Dynamic tables and materialized views
    6. 6.
      “Materialized views accelerate repeated queries against a single base table.”
      ↩︎ Dynamic tables and materialized views
    7. 7.
      “Data accessed through materialized views is always current, regardless of the amount of DML that has been performed on the base table.”
      ↩︎ Dynamic tables and materialized views
    8. 8.
      “Depending on how you configure it, the function can return either scalar results or tabular results.”
      ↩︎ UDFs, Snowflake Scripting and Notebooks
    9. 9.
      “A simple block only requires the keywords BEGIN and END.”
      ↩︎ UDFs, Snowflake Scripting and Notebooks
      “The keyword BEGIN that starts a block is different from the keyword BEGIN that starts a transaction.”
      ↩︎ Checkpoint
    10. 10.
      “The Snowflake SQL API is a REST API that you can use to access and update data in a Snowflake database.”
      ↩︎ Driving pipelines through the SQL API
      “The following commands and statements are supported only within a request that specifies multiple statements:”
      ↩︎ Driving pipelines through the SQL API
    11. 11.
      “Use this endpoint to cancel the execution of a statement.”
      ↩︎ Driving pipelines through the SQL API
    12. 12.
      “Use OAuth or Key Pair to authenticate with the Snowflake server.”
      ↩︎ Driving pipelines through the SQL API
    13. 13.
      “The execution of the statement exceeded the timeout period. The execution of the statement was cancelled.”
      ↩︎ Driving pipelines through the SQL API
      “An error occurred when executing the statement. Check the error code and error message for details.”
      ↩︎ Driving pipelines through the SQL API
      “a statement is executed and the results are returned if the execution is completed in 45 seconds.”
      ↩︎ Checkpoint

    Ready to test yourself?

    Practise the 16 questions on this subdomain.

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