mirror of
https://github.com/supabase/supabase.git
synced 2026-10-05 09:25:06 +03:00
docs(pipelines): add Snowflake materialization examples (#50571)
## I have read the [CONTRIBUTING.md](https://github.com/supabase/supabase/blob/master/CONTRIBUTING.md) file. YES ## What kind of change does this PR introduce? Documentation update. ## What is the current behavior? The Snowflake destination guide describes its append-only change history, but does not include SQL examples for querying current state or maintaining a materialized result. ## What is the new behavior? Add a "Query and materialize current state" section with: - A query and reusable view that select the latest event per identity before filtering deletes. - An incremental dynamic-table example with a configurable freshness target. - Guidance on stable keys, permissions, change tracking, refresh costs, and recovery after table resets or schema changes. - Links to official Snowflake documentation, including the streams-and-tasks alternative. <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Documentation** * Expanded Snowflake replication guidance for deriving current state from append-only change history. * Added examples for identity selection, `QUALIFY`-based filtering, reusable views, dynamic tables, streams, and tasks. * Documented considerations for mutable identity columns, delete handling, change tracking permissions, refresh settings, target lag, and DDL effects. * Clarified that change tracking must be enabled before altering managed objects. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
1 parent
45a80b5685
commit
a536a8fdd0
1 file changed
+137
-2
@@ -154,10 +154,145 @@ Snowflake tables are an event history, not a current-state replica:
|
||||
- A delete appends the complete old row for `REPLICA IDENTITY FULL`. For a primary-key or `USING INDEX` identity, it appends only the identity columns and sets all other source columns to `NULL`.
|
||||
- A source `TRUNCATE` truncates the Snowflake table, resets its streaming state, and does not append a truncate event.
|
||||
|
||||
To derive current state, group by a stable source identity and select the row with the latest `_cdc_sequence_number`. Exclude identities whose latest operation is `delete`. The sequence number is used for ordering and checkpointing. It is not a globally unique event ID. Pipelines provides at-least-once delivery, so consumers must tolerate duplicates. Snowpipe committed offsets suppress routine replay but do not change this guarantee.
|
||||
To derive current state, group by a stable source identity and select the row with the latest `_cdc_sequence_number`. Exclude identities whose latest operation is `delete`. See [Query and materialize current state](#query-and-materialize-current-state) for SQL examples.
|
||||
|
||||
The sequence number is used for ordering and checkpointing. It is not a globally unique event ID. Pipelines provides at-least-once delivery, so consumers must tolerate duplicates. Snowpipe committed offsets suppress routine replay but do not change this guarantee.
|
||||
|
||||
Resetting a table drops and recreates its Snowflake table and managed streaming state. This erases its history. Removing a table from the Postgres publication stops new changes after the pipeline restarts. The existing Snowflake table remains.
|
||||
|
||||
## Query and materialize current state
|
||||
|
||||
Use the replicated change history to build a current-state dataset for reports and analytics. Pipelines maintains the history table. You create and maintain the queries, views, or dynamic tables that read it.
|
||||
|
||||
| Approach | When to use it | Tradeoff |
|
||||
| -------------------------------------------------- | --------------------------------------------------------- | -------------------------------------------------------------------------------------- |
|
||||
| [Query or view](#query-current-state) | Read current state from the changes already in Snowflake. | Computes the result when queried, so query cost can grow with the history. |
|
||||
| [Dynamic table](#materialize-with-a-dynamic-table) | Store current state for repeated analytics queries. | Uses compute and storage to maintain the result, with a configurable freshness target. |
|
||||
| [Streams and tasks](#use-streams-and-tasks) | Control how and when a separate table is updated. | Requires your own merge, initialization, and recovery logic. |
|
||||
|
||||
### Before you start
|
||||
|
||||
The examples use `public.orders`, replicated to `PIPELINES_DB.REPLICATED.PUBLIC_ORDERS`, with source columns `id` and `status`. Replace these names with your own. Wait for the table's initial sync to finish before treating the result as a complete replica.
|
||||
|
||||
Choose a unique, non-null identity that stays the same when a row is updated. The examples use `id`. For a composite key, include every key column in `partition by`, such as `partition by "tenant_id", "id"`. Include those columns in the publication and in delete events. `REPLICA IDENTITY FULL` alone does not make rows unique.
|
||||
|
||||
<Admonition type="caution">
|
||||
|
||||
Changing an identity column can leave the old identity in these results. Pipelines appends the new row for an update without a delete for the previous identity. Use an immutable key for this pattern.
|
||||
|
||||
</Admonition>
|
||||
|
||||
Use a separate analytics role and warehouse, with a schema outside the Pipelines-managed `REPLICATED` schema for derived objects. The examples use `ANALYTICS_ROLE`, `ANALYTICS_WH`, and `PIPELINES_DB.ANALYTICS`. Ask your Snowflake administrator to prepare these resources and grant the analytics role:
|
||||
|
||||
- `USAGE` on the warehouse, database, and both schemas.
|
||||
- `SELECT` on the replicated table.
|
||||
- `CREATE VIEW` on the analytics schema to create a view, or `CREATE DYNAMIC TABLE` to create a dynamic table.
|
||||
|
||||
The role must be available to the Snowflake user running the examples. Keep ownership of the replicated table with `PIPELINES_ROLE`. See Snowflake's [dynamic table access control](https://docs.snowflake.com/en/user-guide/dynamic-tables/privileges) for the full privilege requirements.
|
||||
|
||||
### Query current state
|
||||
|
||||
Run these statements in a Snowflake SQL worksheet with your analytics role:
|
||||
|
||||
```sql
|
||||
use role ANALYTICS_ROLE;
|
||||
use warehouse ANALYTICS_WH;
|
||||
|
||||
select "id", "status"
|
||||
from PIPELINES_DB.REPLICATED.PUBLIC_ORDERS
|
||||
qualify row_number() over (
|
||||
partition by "id" order by "_cdc_sequence_number" desc
|
||||
) = 1
|
||||
and "_cdc_operation" != 'delete';
|
||||
```
|
||||
|
||||
The result contains one row per identity whose latest operation is not `delete`. Ordering by the fixed-width sequence string selects the latest change. Repeated copies of the same event produce one result row. Keep the double quotes around source and metadata column names because Pipelines creates them as case-sensitive identifiers.
|
||||
|
||||
Keep the delete condition in `qualify`. A `where "_cdc_operation" != 'delete'` condition would remove delete events before ranking and could bring back an older row. Snowflake's [`QUALIFY` reference](https://docs.snowflake.com/en/sql-reference/constructs/qualify) explains this evaluation order.
|
||||
|
||||
To reuse the query from an analytics tool, save it as a view:
|
||||
|
||||
```sql
|
||||
create view PIPELINES_DB.ANALYTICS.ORDERS_CURRENT_VIEW as
|
||||
select "id", "status"
|
||||
from PIPELINES_DB.REPLICATED.PUBLIC_ORDERS
|
||||
qualify row_number() over (
|
||||
partition by "id" order by "_cdc_sequence_number" desc
|
||||
) = 1
|
||||
and "_cdc_operation" != 'delete';
|
||||
```
|
||||
|
||||
A regular view stores the query definition, not a separate copy of its results. Each read derives current state from the history available to that query. See Snowflake's [comparison of views and dynamic tables](https://docs.snowflake.com/en/user-guide/overview-view-mview-dts).
|
||||
|
||||
### Materialize with a dynamic table
|
||||
|
||||
A dynamic table stores the query result and refreshes it as the replicated history changes. Use it when you want to query a maintained current-state dataset without defining a scheduled merge task.
|
||||
|
||||
1. Ask the owner of the replicated table to enable change tracking in Snowflake. This is a table setting, not a change to the replicated columns or data. Run as `PIPELINES_ROLE`, or another role that inherits ownership:
|
||||
|
||||
```sql
|
||||
alter table PIPELINES_DB.REPLICATED.PUBLIC_ORDERS
|
||||
set change_tracking = true;
|
||||
```
|
||||
|
||||
The analytics role does not own the replicated table, so it cannot enable change tracking automatically when creating the dynamic table. See Snowflake's [change tracking requirements](https://docs.snowflake.com/en/user-guide/dynamic-tables/troubleshoot-creation#change-tracking-not-enabled-on-base-tables).
|
||||
|
||||
2. Switch to the analytics role and create the dynamic table:
|
||||
|
||||
```sql
|
||||
use role ANALYTICS_ROLE;
|
||||
use warehouse ANALYTICS_WH;
|
||||
|
||||
create dynamic table PIPELINES_DB.ANALYTICS.ORDERS_CURRENT
|
||||
target_lag = '5 minutes'
|
||||
warehouse = ANALYTICS_WH
|
||||
refresh_mode = incremental
|
||||
initialize = on_create
|
||||
as
|
||||
select "id", "status"
|
||||
from PIPELINES_DB.REPLICATED.PUBLIC_ORDERS
|
||||
qualify row_number() over (
|
||||
partition by "id" order by "_cdc_sequence_number" desc
|
||||
) = 1
|
||||
and "_cdc_operation" != 'delete';
|
||||
```
|
||||
|
||||
`initialize = on_create` populates the dynamic table before creation finishes. Explicit `refresh_mode = incremental` makes creation fail if your adapted query cannot refresh incrementally, instead of choosing a full refresh through `AUTO`. See Snowflake's [refresh modes](https://docs.snowflake.com/en/user-guide/dynamic-tables/refresh-modes) and [`CREATE DYNAMIC TABLE` reference](https://docs.snowflake.com/en/sql-reference/sql/create-dynamic-table).
|
||||
|
||||
3. Check the refresh mode and read the materialized rows:
|
||||
|
||||
```sql
|
||||
show dynamic tables like 'ORDERS_CURRENT'
|
||||
in schema PIPELINES_DB.ANALYTICS;
|
||||
|
||||
select "id", "status"
|
||||
from PIPELINES_DB.ANALYTICS.ORDERS_CURRENT;
|
||||
```
|
||||
|
||||
Confirm that `refresh_mode` is `INCREMENTAL` and scheduling is running. Use [Snowflake's refresh monitoring](https://docs.snowflake.com/en/user-guide/dynamic-tables/monitoring) to check the last successful refresh and any errors. After an insert, update, or delete reaches the replicated table, the next successful refresh reflects it in `ORDERS_CURRENT`.
|
||||
|
||||
The five-minute `target_lag` is an example freshness target relative to the history in Snowflake. It is not a fixed refresh schedule or an end-to-end latency guarantee from Postgres. Pipeline replication lag and dynamic-table refresh lag both affect freshness. See Snowflake's [target lag guide](https://docs.snowflake.com/en/user-guide/dynamic-tables/target-lag).
|
||||
|
||||
Dynamic-table refreshes consume warehouse compute, and the materialized results consume storage. These costs are additional to ingestion and querying. Start with a freshness target that meets your reporting needs and measure a representative workload. A dedicated warehouse helps isolate refresh costs. See Snowflake's [dynamic table cost guide](https://docs.snowflake.com/en/user-guide/dynamic-tables/cost).
|
||||
|
||||
### Maintain derived objects
|
||||
|
||||
Pipelines maintains the replicated history table, but does not update your view or dynamic-table definitions.
|
||||
|
||||
| Change | What to do |
|
||||
| ---------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
|
||||
| Source `TRUNCATE` | A direct query or view reads the truncated history. Check that the dynamic table completes a refresh before relying on its contents. |
|
||||
| Pipeline table reset | Wait for the new initial sync. Reapply table-specific read grants and change tracking to the recreated history table. Check dependent objects and recreate the dynamic table if it cannot refresh. |
|
||||
| Added, renamed, or dropped source column | Review the explicit column list. Add new columns to your definition when needed. Update or recreate derived objects that reference renamed or dropped columns. |
|
||||
|
||||
Recreating a dynamic table initializes its contents again and uses compute. See Snowflake's [dynamic table modification guide](https://docs.snowflake.com/en/user-guide/dynamic-tables/modify) for changes that require reinitialization.
|
||||
|
||||
### Use streams and tasks
|
||||
|
||||
Snowflake [streams and tasks](https://docs.snowflake.com/en/user-guide/data-pipelines-intro) can maintain a separate table with scheduled `MERGE` statements. Use this option when you need control over the update procedure or schedule. Snowflake's [SCD Type 1 examples](https://docs.snowflake.com/en/user-guide/dynamic-tables/migrate-streams-tasks#scd-type-1-upsert) compare this approach with dynamic tables.
|
||||
|
||||
Adapt the merge to Pipelines' `"_cdc_operation"` and `"_cdc_sequence_number"` columns. A stream on the history table sees appended rows, including rows representing source updates and deletes. Your job must interpret those operations, load existing history, tolerate replay, and rebuild current state after a source truncate or pipeline table reset.
|
||||
|
||||
## Source table requirements
|
||||
|
||||
Required `REPLICA IDENTITY` depends on the operations enabled in the Postgres publication:
|
||||
@@ -218,7 +353,7 @@ Unsupported or limited changes:
|
||||
- Changes to nullability or existing column defaults are ignored.
|
||||
- Initial table creation can copy compatible literal defaults. Added columns can copy string, numeric, or boolean literal defaults. Other defaults are omitted.
|
||||
|
||||
Snowflake DDL changes existing history. Adding a column with a default can populate older rows. Renaming a column changes the historical schema. Dropping a column removes it from old events. Snowflake DDL is not transactional, so an interrupted multi-column change can leave a partially applied schema. Do not alter managed destination objects manually. If the pipeline remains failed after a restart, [contact support](/dashboard/support/new).
|
||||
Snowflake DDL changes existing history. Adding a column with a default can populate older rows. Renaming a column changes the historical schema. Dropping a column removes it from old events. Snowflake DDL is not transactional, so an interrupted multi-column change can leave a partially applied schema. Apart from [enabling change tracking](#materialize-with-a-dynamic-table), do not alter managed destination objects manually. If the pipeline remains failed after a restart, [contact support](/dashboard/support/new).
|
||||
|
||||
## Troubleshooting
|
||||
|
||||
|
||||
Reference in new issue
Block a user