Snowflake Dynamic Tables for Streaming Pipelines

By Sophia Bennett • • 9 min read

Building streaming data pipelines has traditionally meant maintaining two separate systems: a stream processor like Apache Flink or Kafka Streams for real-time transformations, and a warehouse like Snowflake for batch analytics. This dual architecture doubles operational complexity, introduces consistency challenges between real-time and batch views, and demands engineers who understand both paradigms.

Snowflake Dynamic Tables collapse this divide. Instead of writing imperative streaming code that manages state, watermarks, and checkpoints, you write a declarative SQL query that defines what you want the result to look like. Snowflake handles the when and how — automatically scheduling incremental refreshes to keep the table within a specified staleness target. The result is a streaming pipeline built entirely in SQL, managed entirely by Snowflake, with no external orchestrator required.

This article explores how Dynamic Tables work under the hood, how to design multi-layer streaming pipelines with them, and where they fit alongside dedicated stream processing frameworks like those covered in our guide on stream-stream joins in Flink and Kafka Streams.

Understanding Dynamic Tables

A Dynamic Table is a Snowflake table whose contents are defined by a SQL query. Unlike a view, which re-executes the query on every read, a Dynamic Table materializes the results and keeps them fresh through automatic, incremental refreshes. Unlike a traditional materialized view, a Dynamic Table supports the full breadth of Snowflake SQL — joins, window functions, CTEs, subqueries, and even calls to user-defined functions.

The core concept is simple: you declare the transformation, and Snowflake manages the execution. Here is a basic Dynamic Table that aggregates streaming events into hourly metrics:

CREATE OR REPLACE DYNAMIC TABLE hourly_metrics
  TARGET_LAG = '5 minutes'
  WAREHOUSE = transform_wh
AS
SELECT
    DATE_TRUNC('hour', event_timestamp) AS hour,
    event_type,
    COUNT(*) AS event_count,
    COUNT(DISTINCT user_id) AS unique_users,
    AVG(duration_ms) AS avg_duration_ms,
    PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY duration_ms) AS p95_duration_ms
FROM raw_events
GROUP BY 1, 2;

The TARGET_LAG parameter is the key differentiator. It tells Snowflake the maximum acceptable staleness for this table. With a target lag of 5 minutes, Snowflake guarantees that the data in hourly_metrics will never be more than 5 minutes behind the data in raw_events. Snowflake automatically determines the refresh frequency needed to meet this target, and it uses incremental processing — scanning only the changed data — to minimize compute costs.

Incremental Refresh Mechanics

Dynamic Tables use Snowflake's change tracking infrastructure to identify which rows in the source tables have been inserted, updated, or deleted since the last refresh. For queries that Snowflake can incrementally maintain — aggregations, filters, projections, and certain types of joins — it processes only the changed data and merges the results into the existing materialization.

When a query cannot be incrementally maintained (for example, queries with certain window functions or non-deterministic expressions), Snowflake falls back to a full refresh. You can check whether your Dynamic Table qualifies for incremental refresh by examining its refresh mode:

-- Check refresh mode for a Dynamic Table
SHOW DYNAMIC TABLES LIKE 'hourly_metrics';
-- Look at the 'scheduling_state' and 'refresh_mode' columns
-- INCREMENTAL = only changed data processed
-- FULL = entire query re-executed on each refresh

Building Multi-Layer Pipelines

The real power of Dynamic Tables emerges when you chain them together into multi-layer transformation pipelines. Each layer is a Dynamic Table that references upstream Dynamic Tables, and Snowflake automatically manages the dependency graph and refresh ordering.

Consider a typical medallion architecture for processing clickstream data:

Bronze Layer: Raw Ingestion

The bronze layer lands raw data from Snowpipe Streaming into a staging table. This is a regular table, not a Dynamic Table, since it is the entry point for streaming ingestion:

-- Bronze: raw events from Snowpipe Streaming
CREATE TABLE raw_clickstream (
    event_id STRING,
    user_id STRING,
    session_id STRING,
    page_url STRING,
    event_type STRING,   -- 'pageview', 'click', 'scroll'
    event_payload VARIANT,
    event_timestamp TIMESTAMP_NTZ,
    ingested_at TIMESTAMP_NTZ DEFAULT CURRENT_TIMESTAMP()
)
CHANGE_TRACKING = TRUE;  -- Required for downstream Dynamic Tables

Note the CHANGE_TRACKING = TRUE setting. Source tables for Dynamic Tables must have change tracking enabled so that Snowflake can identify new and modified rows for incremental processing.

Silver Layer: Cleaned and Enriched

The silver layer applies data quality rules, deduplicates events, and enriches with dimension data:

CREATE OR REPLACE DYNAMIC TABLE silver_clickstream
  TARGET_LAG = '2 minutes'
  WAREHOUSE = transform_wh
AS
SELECT
    rc.event_id,
    rc.user_id,
    rc.session_id,
    rc.page_url,
    rc.event_type,
    rc.event_timestamp,
    -- Extract structured fields from semi-structured payload
    rc.event_payload:element_id::STRING AS element_id,
    rc.event_payload:viewport_width::INT AS viewport_width,
    rc.event_payload:scroll_depth::FLOAT AS scroll_depth,
    -- Enrich with user dimensions
    u.account_tier,
    u.signup_date,
    u.country_code,
    -- Derive session sequence
    ROW_NUMBER() OVER (
        PARTITION BY rc.session_id
        ORDER BY rc.event_timestamp
    ) AS event_sequence
FROM raw_clickstream rc
LEFT JOIN dim_users u ON rc.user_id = u.user_id
WHERE rc.event_timestamp >= DATEADD('day', -7, CURRENT_TIMESTAMP())
  AND rc.event_id IS NOT NULL
QUALIFY ROW_NUMBER() OVER (
    PARTITION BY rc.event_id ORDER BY rc.ingested_at DESC
) = 1;  -- Deduplicate by event_id

Gold Layer: Business Aggregates

The gold layer produces business-ready aggregations consumed by dashboards and applications:

CREATE OR REPLACE DYNAMIC TABLE gold_engagement_metrics
  TARGET_LAG = '5 minutes'
  WAREHOUSE = analytics_wh
AS
SELECT
    DATE_TRUNC('hour', event_timestamp) AS metric_hour,
    country_code,
    account_tier,
    COUNT(DISTINCT session_id) AS total_sessions,
    COUNT(DISTINCT user_id) AS unique_users,
    COUNT(CASE WHEN event_type = 'pageview' THEN 1 END) AS pageviews,
    COUNT(CASE WHEN event_type = 'click' THEN 1 END) AS clicks,
    AVG(scroll_depth) AS avg_scroll_depth,
    -- Bounce rate: sessions with only one event
    COUNT(DISTINCT CASE
        WHEN event_sequence = 1
        AND session_id NOT IN (
            SELECT session_id FROM silver_clickstream
            WHERE event_sequence > 1
        )
        THEN session_id
    END)::FLOAT / NULLIF(COUNT(DISTINCT session_id), 0) AS bounce_rate
FROM silver_clickstream
GROUP BY 1, 2, 3;

In this pipeline, gold_engagement_metrics depends on silver_clickstream, which depends on raw_clickstream. Snowflake recognizes this dependency chain and refreshes them in order. The gold table has a 5-minute target lag, and the silver table has a 2-minute target lag, so the end-to-end latency from raw ingestion to business metrics is at most 7 minutes.

Target Lag Strategies

Choosing the right TARGET_LAG for each Dynamic Table involves balancing freshness requirements against compute costs. Shorter lag means more frequent refreshes, which consume more warehouse credits.

Use Case Recommended Lag Rationale
Real-time dashboards 1-5 minutes Users expect near-live data; refresh cost justified by visibility
Operational reporting 10-30 minutes Business decisions tolerate moderate delay; lower compute cost
Daily aggregations 1-6 hours Consumers check once or twice daily; minimize warehouse usage
Intermediate layers DOWNSTREAM Let downstream consumers drive the refresh cadence

The DOWNSTREAM option is particularly useful for intermediate transformation layers. Instead of specifying a fixed lag, DOWNSTREAM tells Snowflake to refresh this table on demand when a downstream Dynamic Table needs fresh data. This avoids unnecessary refreshes of intermediate tables that no one queries directly.

-- Intermediate table refreshes only when downstream tables need it
CREATE OR REPLACE DYNAMIC TABLE silver_orders_enriched
  TARGET_LAG = DOWNSTREAM
  WAREHOUSE = transform_wh
AS
SELECT
    o.*,
    c.customer_segment,
    p.product_category
FROM raw_orders o
JOIN dim_customers c ON o.customer_id = c.customer_id
JOIN dim_products p ON o.product_id = p.product_id;

Snowpipe Streaming Integration

Dynamic Tables reach their full potential when paired with Snowpipe Streaming, which provides low-latency ingestion from Kafka, Kinesis, or custom producers directly into Snowflake tables. The combination creates an end-to-end streaming pipeline entirely within Snowflake.

The ingestion flow works as follows: your application produces events to Kafka, the Snowflake Kafka Connector with Snowpipe Streaming mode writes rows into a landing table with sub-second latency, and the Dynamic Table chain transforms and aggregates the data within the specified target lag.

Key configuration for the Kafka connector in streaming mode:

# Snowflake Kafka Connector - Snowpipe Streaming config
connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector
snowflake.ingestion.method=SNOWPIPE_STREAMING
snowflake.enable.schematization=true
snowflake.role.name=STREAMING_ROLE
snowflake.database.name=ANALYTICS
snowflake.schema.name=RAW
buffer.count.records=10000
buffer.flush.time=10
buffer.size.bytes=5000000

With snowflake.ingestion.method=SNOWPIPE_STREAMING, the connector bypasses the traditional stage-and-copy pattern and writes directly to Snowflake's rowstore, making data available for Dynamic Table consumption within seconds rather than minutes.

Monitoring and Troubleshooting

Dynamic Tables provide built-in observability through the INFORMATION_SCHEMA and ACCOUNT_USAGE schemas. Monitoring refresh performance and costs is essential for maintaining efficient pipelines.

Refresh History

Query the refresh history to understand how your Dynamic Tables are performing:

-- Recent refresh history with performance metrics
SELECT
    name AS dynamic_table_name,
    refresh_trigger,           -- SCHEDULED or ON_DEMAND
    refresh_action,            -- INCREMENTAL or FULL
    state,                     -- SUCCEEDED, FAILED, CANCELLED
    data_timestamp,
    refresh_start_time,
    refresh_end_time,
    DATEDIFF('second', refresh_start_time, refresh_end_time) AS refresh_seconds,
    statistics:numInsertedRows::INT AS rows_inserted,
    statistics:numDeletedRows::INT AS rows_deleted,
    statistics:numUpdatedRows::INT AS rows_updated
FROM TABLE(INFORMATION_SCHEMA.DYNAMIC_TABLE_REFRESH_HISTORY(
    NAME => 'ANALYTICS.PUBLIC.GOLD_ENGAGEMENT_METRICS'
))
ORDER BY refresh_start_time DESC
LIMIT 20;

Lag Monitoring

Track whether your Dynamic Tables are meeting their target lag commitments:

-- Current lag for all Dynamic Tables
SELECT
    name,
    target_lag,
    TIMEDIFF('second', data_timestamp, CURRENT_TIMESTAMP()) AS current_lag_seconds,
    scheduling_state,
    refresh_mode
FROM TABLE(INFORMATION_SCHEMA.DYNAMIC_TABLES())
WHERE current_lag_seconds > SPLIT_PART(target_lag, ' ', 1)::INT * 60
ORDER BY current_lag_seconds DESC;

If a Dynamic Table consistently exceeds its target lag, common causes include an undersized warehouse, expensive full refreshes (indicating the query is not eligible for incremental mode), or upstream tables refreshing too slowly. Increasing the warehouse size or simplifying the query to qualify for incremental processing usually resolves lag issues.

Dynamic Tables vs. Other Approaches

Dynamic Tables are not the only way to build transformation pipelines in Snowflake. Understanding when to use them versus alternatives helps you make better architectural decisions.

Dynamic Tables vs. Streams and Tasks

Before Dynamic Tables, the standard Snowflake pattern for incremental processing used Streams (change data capture) and Tasks (scheduled execution). This approach works but requires you to write imperative MERGE statements, manage task schedules, handle failure recovery, and maintain the dependency ordering yourself.

Dynamic Tables abstract all of this away. You write the declarative query, and Snowflake manages the incremental logic, scheduling, and dependency resolution. The trade-off is less control: you cannot customize the merge logic, and you cannot pause or reorder individual refresh steps as easily as you can with Tasks.

Dynamic Tables vs. dbt

dbt and Dynamic Tables solve overlapping problems but from different angles. dbt provides version-controlled transformations, testing, documentation, and lineage across any warehouse. Dynamic Tables provide automated, low-latency incremental refresh native to Snowflake. Many teams use both: dbt for defining and testing the transformation logic, and Dynamic Tables as the materialization strategy for models that need near-real-time refresh.

Dynamic Tables vs. External Stream Processors

For use cases requiring sub-second latency, complex event processing, or integration with non-Snowflake systems, dedicated stream processors like Apache Flink remain the better choice. Dynamic Tables excel in the analytical streaming tier — where data from various sources needs to be joined, aggregated, and made available for querying with minute-level freshness. The comparison is similar to how ClickHouse materialized views handle real-time analytical workloads differently than general-purpose stream processors.

Production Best Practices

After operating Dynamic Table pipelines across several production environments, a set of patterns has emerged that consistently improve reliability and cost efficiency.

Warehouse Sizing and Isolation

Assign dedicated warehouses to Dynamic Table refresh workloads. Sharing a warehouse with ad-hoc queries causes unpredictable refresh latency when the warehouse is busy. Start with an X-Small warehouse and scale up based on refresh duration metrics — Snowflake's per-second billing makes right-sizing low-risk.

Change Tracking Overhead

Enabling CHANGE_TRACKING on source tables adds metadata overhead to every DML operation. For high-volume ingestion tables receiving millions of rows per hour, test the write performance impact before enabling change tracking in production. The overhead is typically under 5% but can be higher for tables with many columns.

Handling Late-Arriving Data

Dynamic Tables process data based on when it appears in the source table, not based on event timestamps. If your pipeline receives late-arriving data — events that arrive hours after they occurred — the Dynamic Table will process them on their next refresh cycle. For analytical accuracy, filter on event timestamps and handle late arrivals through change data feed patterns that reprocess affected time windows.

Cost Control

Monitor warehouse credit consumption per Dynamic Table using the refresh history views. Set resource monitors with alerts and suspension policies to prevent runaway costs from unexpected data volume spikes. Consider using TARGET_LAG = DOWNSTREAM for intermediate tables to avoid unnecessary refreshes.

Dynamic Tables represent Snowflake's clearest statement that the boundary between batch and streaming analytics is artificial. By making the refresh interval a tunable parameter on a declarative SQL definition, they let you dial the latency knob from near-real-time to daily without changing your transformation logic.

The practical impact is significant. Teams that previously maintained separate Flink jobs and dbt models for the same business logic can consolidate into a single Dynamic Table pipeline, reducing the operational surface area and eliminating consistency issues between real-time and batch views. For organizations already invested in Snowflake, Dynamic Tables are the most natural path to streaming analytics without the overhead of managing a separate stream processing infrastructure.

Start with a single use case — a dashboard that currently refreshes hourly but needs minute-level freshness — and convert its transformation chain to Dynamic Tables. Measure the cost, monitor the lag, and expand from there as the pattern proves itself in your environment.