May 16, 2024

Snowflake Dynamic Tables Are GA: Replacing Streams and Tasks Pipelines

Snowflake announced at the end of April that Dynamic Tables are generally available. I've been trying them out in preview on a sample bordereaux model, and they look like a strong replacement for many Streams + Tasks chains. This post explains what they are, how they compare, and where Streams and Tasks still make sense.

The examples use a delegated authority model: coverholders submit premium bordereaux under binding authority contract sections, and we build a clean premium transaction table and a GWP mart by coverholder and year of account.

The old way: Streams + Tasks

I've used Streams and Tasks for incremental bordereaux loads for a while. A typical pipeline looks like this:

 1: -- 1. Capture changes on the raw premium bordereau table
 2: CREATE OR REPLACE STREAM raw.premium_bdx_stream ON TABLE raw.premium_bdx;
 3: 
 4: -- 2. A task that runs every 15 minutes when new rows have arrived
 5: CREATE OR REPLACE TASK transform.load_premium_transaction
 6:   WAREHOUSE = transform_wh
 7:   SCHEDULE  = '15 MINUTE'
 8: WHEN SYSTEM$STREAM_HAS_DATA('raw.premium_bdx_stream')
 9: AS
10: MERGE INTO clean.premium_transaction t
11: USING (
12:   SELECT * FROM raw.premium_bdx_stream
13:   QUALIFY ROW_NUMBER() OVER (
14:             PARTITION BY umr, section_no, certificate_ref, transaction_seq
15:             ORDER BY submitted_at DESC) = 1
16: ) s
17: ON  t.umr = s.umr AND t.section_no = s.section_no
18: AND t.certificate_ref = s.certificate_ref AND t.transaction_seq = s.transaction_seq
19: WHEN MATCHED AND s.METADATA$ACTION = 'DELETE' AND NOT s.METADATA$ISUPDATE THEN DELETE
20: WHEN MATCHED THEN UPDATE SET
21:      t.gross_premium = s.gross_premium,
22:      t.commission    = s.commission,
23:      t.submitted_at  = s.submitted_at
24: WHEN NOT MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
25:   INSERT (umr, section_no, certificate_ref, transaction_seq, transaction_type,
26:           bordereau_month, original_ccy, gross_premium, commission, submitted_at)
27:   VALUES (s.umr, s.section_no, s.certificate_ref, s.transaction_seq, s.transaction_type,
28:           s.bordereau_month, s.original_ccy, s.gross_premium, s.commission, s.submitted_at);
29: 
30: ALTER TASK transform.load_premium_transaction RESUME;

It works, but you're writing the how by hand: the stream offsets, the MERGE logic, the handling of deletes and updates, the task graph, and resuming child tasks in the right order. Multiply that by risk, premium and claims bordereaux, plus every dimension, and the pipeline turns into a maintenance burden.

The new way: Dynamic Tables

A Dynamic Table is defined by a query and a target lag. Snowflake works out how to keep the result up to date, incrementally where it can.

 1: CREATE OR REPLACE DYNAMIC TABLE clean.premium_transaction
 2:   TARGET_LAG = '15 minutes'
 3:   WAREHOUSE  = transform_wh
 4: AS
 5: SELECT umr, section_no, certificate_ref, transaction_seq, transaction_type,
 6:        bordereau_month, original_ccy, gross_premium, commission, submitted_at
 7: FROM raw.premium_bdx
 8: QUALIFY ROW_NUMBER() OVER (
 9:           PARTITION BY umr, section_no, certificate_ref, transaction_seq
10:           ORDER BY submitted_at DESC) = 1;

That's the whole thing. No stream, no MERGE, no task. You describe the result you want and how fresh it needs to be. When a coverholder resubmits a corrected bordereau, the latest submission simply wins.

Chaining Dynamic Tables

Dynamic Tables can read from other Dynamic Tables, and Snowflake manages the dependency graph. For intermediate tables, use TARGET_LAG = DOWNSTREAM, which means "refresh only when something downstream needs me".

 1: -- Contract sections from the binding authority register
 2: CREATE OR REPLACE DYNAMIC TABLE clean.contract_section
 3:   TARGET_LAG = DOWNSTREAM
 4:   WAREHOUSE  = transform_wh
 5: AS
 6: SELECT umr, section_no, coverholder_id, broker_code,
 7:        class_of_business, year_of_account, gpi_limit_gbp
 8: FROM raw.binding_authority_section
 9: WHERE status <> 'NTU';   -- not taken up
10: 
11: -- Premium converted to GBP at the month rate of exchange
12: CREATE OR REPLACE DYNAMIC TABLE clean.premium_transaction_gbp
13:   TARGET_LAG = DOWNSTREAM
14:   WAREHOUSE  = transform_wh
15: AS
16: SELECT p.*,
17:        p.gross_premium * r.rate_to_gbp AS gross_premium_gbp,
18:        p.commission    * r.rate_to_gbp AS commission_gbp
19: FROM clean.premium_transaction p
20: JOIN raw.month_roe r
21:   ON r.currency = p.original_ccy AND r.roe_month = p.bordereau_month;
22: 
23: -- The mart that oversight dashboards read
24: CREATE OR REPLACE DYNAMIC TABLE mart.gwp_by_coverholder_yoa
25:   TARGET_LAG = '1 hour'
26:   WAREHOUSE  = transform_wh
27: AS
28: SELECT cs.coverholder_id,
29:        cs.broker_code,
30:        cs.class_of_business,
31:        cs.year_of_account,
32:        p.bordereau_month,
33:        SUM(p.gross_premium_gbp)                       AS gwp_gbp,
34:        SUM(p.commission_gbp)                          AS commission_gbp,
35:        SUM(p.gross_premium_gbp - p.commission_gbp)    AS net_premium_gbp,
36:        COUNT(DISTINCT p.certificate_ref)              AS certificates
37: FROM clean.premium_transaction_gbp p
38: JOIN clean.contract_section cs
39:   ON cs.umr = p.umr AND cs.section_no = p.section_no
40: GROUP BY 1, 2, 3, 4, 5;

Only the final table has a real time target. The intermediate ones refresh as needed to meet it. In Snowsight you can see the whole graph, with refresh history and lag, under the Dynamic Table's Graph tab, which is also a nice thing to show an auditor asking how a number on the dashboard was built.

Incremental vs full refresh

This is the part to understand properly. Each Dynamic Table has a refresh mode: INCREMENTAL, FULL, or AUTO (the default, which picks one when the table is created). Incremental refresh only processes changes and is what makes Dynamic Tables cheap. But not every query can be maintained incrementally. Some constructs, certain non-deterministic functions for example, force a full refresh.

Check what you got:

1: SHOW DYNAMIC TABLES LIKE 'gwp_by_coverholder_yoa' IN SCHEMA mart;
2: -- look at the refresh_mode and refresh_mode_reason columns

If a large table ends up on FULL, every refresh recomputes everything, which can cost more than the Streams + Tasks pipeline you replaced. It's worth setting REFRESH_MODE = INCREMENTAL explicitly on big tables like the premium and claims transactions, so creation fails if the query isn't supported, rather than quietly falling back.

Monitoring

1: SELECT name, state, refresh_trigger, refresh_start_time, refresh_end_time,
2:        state_message
3: FROM TABLE(INFORMATION_SCHEMA.DYNAMIC_TABLE_REFRESH_HISTORY(
4:        NAME_PREFIX => 'DA_PROD.MART.'))
5: ORDER BY refresh_start_time DESC
6: LIMIT 50;

You can also suspend and resume, or force a refresh after a backfill, such as reloading a coverholder's whole year of account:

1: ALTER DYNAMIC TABLE mart.gwp_by_coverholder_yoa SUSPEND;
2: ALTER DYNAMIC TABLE mart.gwp_by_coverholder_yoa RESUME;
3: ALTER DYNAMIC TABLE mart.gwp_by_coverholder_yoa REFRESH;

Cost tips

  • Target lag drives cost. A 1-minute lag keeps the warehouse busy. Bordereaux arrive monthly, so ask what the business really needs. An hour is usually plenty, even in the busy first week of the month.
  • Use a dedicated, small warehouse for Dynamic Table refreshes, so you can see their cost separately in WAREHOUSE_METERING_HISTORY.
  • Share a target lag across tables that read the same sources, so refreshes line up and the warehouse wakes up less often.

When Streams + Tasks still make sense

  • Side effects: emailing a coverholder when their bordereau fails validation, or writing to several targets from one change set.
  • Procedural logic: anything that needs a stored procedure, branching, or error handling per batch.
  • Exact control over timing, such as "run the month-end close at 18:00 on working day 5". Dynamic Tables are about freshness, not schedules.
  • Append-only audit trails where you need to keep every bordereau submission and change event, not just the current state.

For the usual "clean it, join it, aggregate it" transformations, Dynamic Tables are the better default. Less code, a visible dependency graph, and Snowflake handles the incremental work.