Files
supabase/apps/docs/content/guides/database/replication/pipelines/snowflake.mdx
T

371 lines
27 KiB
Plaintext

---
id: 'snowflake-destination'
title: 'Snowflake destination'
description: 'Configure Snowflake as a Supabase Pipelines destination.'
subtitle: 'Replicate Supabase Postgres changes to Snowflake.'
sidebar_label: 'Snowflake'
---
<$Partial path="pipelines-public-alpha.mdx" />
[Snowflake](https://www.snowflake.com/) is a managed data platform. Supabase Pipelines replicates each Postgres table to Snowflake as an append-only history of changes.
To replicate data to Snowflake:
1. [Prepare a database, schema, role, service user, and key pair](#prepare-snowflake-resources) in Snowflake.
2. [Configure the Snowflake destination](#configure-snowflake-as-a-destination) in the Dashboard.
3. [Query or materialize the replicated data](#query-and-materialize-current-state) in Snowflake.
## Source table requirements
The operations enabled in the Postgres publication determine the required `REPLICA IDENTITY`:
| Published operations | Required replica identity |
| -------------------- | ---------------------------------------------------------------------------------------------------------- |
| `INSERT` only | No row identity is required. |
| `DELETE` | A primary key, an identity index (`USING INDEX`), or full identity (`FULL`). Publish all identity columns. |
| `UPDATE` | `REPLICA IDENTITY FULL`. |
Set full replica identity before publishing updates:
```sql
alter table public.your_table replica identity full;
```
`REPLICA IDENTITY FULL` increases WAL volume but allows Pipelines to construct complete new rows when Postgres omits unchanged out-of-line TOAST values. It applies only to new WAL records. If the retained WAL already contains an incompatible update, change the setting and then restart replication for the affected table.
## Prepare Snowflake resources
Before you create a pipeline, prepare a dedicated Snowflake database and an empty schema for the replicated tables. Pipelines also needs a service user and role that can create and manage those tables. Use unquoted identifiers for the service user and role so Snowflake stores them in uppercase. Pipelines also converts the account and user names to uppercase during authentication.
Run this setup as a Snowflake administrator, changing the example names as needed:
```sql
create role if not exists PIPELINES_ROLE;
create user if not exists PIPELINES_USER
type = service;
grant role PIPELINES_ROLE to user PIPELINES_USER;
alter user PIPELINES_USER set default_role = PIPELINES_ROLE;
create database if not exists PIPELINES_DB;
create schema if not exists PIPELINES_DB.REPLICATED;
grant usage on database PIPELINES_DB to role PIPELINES_ROLE;
grant usage on schema PIPELINES_DB.REPLICATED to role PIPELINES_ROLE;
grant create table on schema PIPELINES_DB.REPLICATED to role PIPELINES_ROLE;
```
In this example, Pipelines creates the destination tables using `PIPELINES_ROLE`. Whatever name you choose, the role must retain ownership so Pipelines can alter, truncate, or drop the tables when required. Do not create the tables in advance under another role.
Snowflake automatically creates a managed default pipe named `<TABLE>-STREAMING` for each table. You don't need to provide a virtual warehouse, stage, or pipe for ingestion. See [Snowpipe Streaming access privileges](https://docs.snowflake.com/en/user-guide/snowpipe-streaming/snowpipe-streaming-access-control) for the required permissions.
To keep ingestion resources separate from downstream workloads, use a different role and warehouse for queries and transformations.
### Keep the SQL and streaming roles aligned
Pipelines connects to Snowflake through separate interfaces for SQL requests and streaming. Both interfaces need to use the same role:
- SQL requests use the optional **Role** under **Advanced settings** in the Dashboard. When **Role** is empty, they use the user's default role.
- Snowpipe Streaming uses the user's `DEFAULT_ROLE`. It does not use the optional **Role** setting.
Grant the permissions above to a dedicated role such as `PIPELINES_ROLE`, then set it as the service user's `DEFAULT_ROLE`. In the Dashboard, either leave **Role** empty or enter the same role explicitly. This keeps SQL validation and streaming aligned.
### Generate a key pair
Pipelines uses an RSA key pair to authenticate with Snowflake. The key must be at least 2048 bits, and Snowflake recommends the PKCS #8 format. Run one of the following commands from the directory where you want to store the key. The command creates a private key named `rsa_key.p8` in that directory.
For an unencrypted private key:
```bash
openssl genrsa 2048 | openssl pkcs8 -topk8 \
-inform PEM -out rsa_key.p8 -nocrypt
```
Or, for a passphrase-protected private key:
```bash
openssl genrsa 2048 | openssl pkcs8 -topk8 -v2 des3 \
-inform PEM -out rsa_key.p8
```
Create a public key from the private key. Snowflake uses the public key to verify connections signed with `rsa_key.p8`:
```bash
openssl rsa -in rsa_key.p8 -pubout -out rsa_key.pub
```
Open `rsa_key.pub` and copy the text between the `BEGIN PUBLIC KEY` and `END PUBLIC KEY` lines. In the following statement, replace the `<public-key-body>` placeholder with the copied text:
```sql
alter user PIPELINES_USER set rsa_public_key = '<public-key-body>';
```
When you configure the destination in the Dashboard, paste or upload the complete `rsa_key.p8` file into **Private key**. Preserve its original PEM header and footer. If the key is passphrase-protected, enter the passphrase into **Private key passphrase**.
Keep `rsa_key.p8` and its passphrase secret. The Dashboard accepts unencrypted PKCS #1 and PKCS #8 keys. It also accepts encrypted PKCS #8 keys with a passphrase, but not encrypted PKCS #1 keys.
See [Snowflake key-pair authentication](https://docs.snowflake.com/en/user-guide/key-pair-auth) to verify the public-key fingerprint and rotate keys with `RSA_PUBLIC_KEY_2`.
### Find the account identifier
Run this query in Snowflake:
```sql
select current_organization_name() || '-' || current_account_name();
```
Use the result as the **Account ID**, for example `MYORG-MYACCOUNT`. The field also accepts legacy one-part account locators, but not full URLs or dotted locator-and-region hostnames. Account IDs can contain up to 63 characters. See [Snowflake account identifiers](https://docs.snowflake.com/en/user-guide/admin-account-identifier) for details.
## Configure Snowflake as a destination
Follow the steps in [Set up Pipelines](/docs/guides/database/replication/pipelines#setup-overview). When prompted to choose a destination, select **Snowflake** and enter the following settings:
| Field | Value |
| -------------------------- | -------------------------------------------------------------------------------------- |
| **Account ID** | `MYORG-MYACCOUNT`, for example; use an organization-account identifier |
| **User** | `PIPELINES_USER`, or your unquoted service user |
| **Database** | `PIPELINES_DB`, or your destination database |
| **Schema** | `REPLICATED`, or your dedicated destination schema |
| **Role** | Leave empty to use the service user's default role, or enter that same role explicitly |
| **Private key** | Paste or upload the complete `rsa_key.p8` private-key PEM file |
| **Private key passphrase** | Only for an encrypted PKCS #8 key |
Enter the database and schema identifiers exactly as stored in Snowflake. Unquoted identifiers are stored in uppercase. Choose an account near the [managed pipeline region](/docs/guides/database/replication/pipelines#region).
Click **Start pipeline**, then complete the validation and cost confirmations.
## How it works
Pipelines uses Snowflake's SQL REST API to validate the database and schema. It also uses the API to create, update, and recreate destination tables and to apply source `TRUNCATE` operations. Pipelines sends initial and ongoing row data through Snowpipe Streaming.
Validation checks authentication, database and schema visibility, and that `QUOTED_IDENTIFIERS_IGNORE_CASE` is `FALSE`. It does not verify that the role can create or own tables or write through Snowpipe Streaming.
### Destination table names
Pipelines combines each Postgres schema and table name into one Snowflake table name. Existing underscores are doubled, the two names are joined with a single underscore, and the result is converted to uppercase:
| Postgres table | Snowflake table |
| ---------------------- | ------------------------ |
| `public.orders` | `PUBLIC_ORDERS` |
| `sales_eu.order_items` | `SALES__EU_ORDER__ITEMS` |
Postgres schema and table names cannot start or end with `_` or contain `"` or `;`. Names that differ only in case map to the same Snowflake name. Use lowercase Postgres names to avoid collisions. Source column names are preserved as quoted identifiers, except for the [reserved metadata names](#append-only-change-history).
### Append-only change history
Each destination table contains the replicated source columns plus two `VARCHAR NOT NULL` metadata columns:
| Column | Meaning |
| ---------------------- | -------------------------------------------------------------------------------------------------------- |
| `_cdc_operation` | Lowercase operation: `insert`, `update`, or `delete`. |
| `_cdc_sequence_number` | Fixed-width hexadecimal commit LSN and transaction ordinal, such as `00000000016b3740/0000000000000002`. |
The metadata names are reserved and can't be used by source columns. Initial-sync rows use `insert` and the shared sequence number `0000000000000000/0000000000000000`.
Snowflake tables contain an event history rather than a current-state replica:
- An insert appends the new row.
- An update appends the complete new row. It does not append a before image.
- A delete appends the complete old row for `REPLICA IDENTITY FULL`. For a primary-key or `USING INDEX` identity, it sends only the identity columns. Other columns can contain destination defaults or `NULL`; do not treat them as the deleted row's original values.
- Truncating the source table also truncates the Snowflake table and resets its streaming state. It does not append a truncate event.
The sequence number orders changes but is not a globally unique event ID. Snowpipe committed offsets suppress routine replay, but consumers must still tolerate [duplicate processing](/docs/guides/database/replication/pipelines-faq#can-data-be-processed-more-than-once).
[Restarting a table](/docs/guides/database/replication/pipelines-monitoring#restarting-tables) drops the Snowflake table and its managed streaming state, which erases the replicated history. A restart cannot recover past events. By contrast, [removing a table from the publication](/docs/guides/database/replication/pipelines#removing-tables-from-replication) leaves its destination history in place.
## Query replicated data [#query-and-materialize-current-state]
Use the replicated change history to build current-state datasets for reporting and analytics. Pipelines maintains the history table, while you maintain the queries, views, or dynamic tables that read from 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 a source table named `public.orders` with the columns `id` and `status`. Pipelines replicates it to `PIPELINES_DB.REPLICATED.PUBLIC_ORDERS`. Replace these names with your own, and wait for the initial sync to finish before treating the results as a complete replica.
Choose a unique, non-null identity that does not change when a row is updated. The examples use `id`. If you use a composite key, include every key column in `partition by`, such as `partition by "tenant_id", "id"`. The publication and delete events must include the same columns. `REPLICA IDENTITY FULL` alone does not make rows unique.
<Admonition type="caution">
If an identity value changes, the old identity can remain in the results. Pipelines appends the updated row without first appending a delete for the previous identity. Use an immutable key for this pattern.
</Admonition>
Create derived objects in a schema outside the Pipelines-managed `REPLICATED` schema, using a separate analytics role and warehouse. 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 query returns one row for each identity whose latest operation is not `delete`. The fixed-width sequence string determines which change is the latest, and repeated copies of the same event collapse into 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 inside `qualify`. Moving it to `where "_cdc_operation" != 'delete'` would remove delete events before the rows are ranked, which could bring back an older version of a deleted 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 rather than a separate copy of the results. Each read derives the current state from the history available at query time. 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 results and refreshes them as the replicated history changes. Use one when you need a maintained current-state dataset without maintaining 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 already in Snowflake. It is neither a fixed refresh schedule nor an end-to-end latency guarantee from Postgres. Both pipeline replication lag and dynamic-table refresh lag affect freshness. See Snowflake's [target lag guide](https://docs.snowflake.com/en/user-guide/dynamic-tables/target-lag).
Refreshing a dynamic table consumes warehouse compute, while its materialized results consume storage. These costs are additional to ingestion and querying. Start with a freshness target that meets your reporting needs, then test its cost and refresh behavior with a representative workload. A dedicated warehouse can help isolate refresh costs. See Snowflake's [dynamic table cost guide](https://docs.snowflake.com/en/user-guide/dynamic-tables/cost).
### Use streams and tasks
Snowflake [streams and tasks](https://docs.snowflake.com/en/user-guide/data-pipelines-intro) can maintain a separate table through scheduled `MERGE` statements. Use them when you need more 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 the `"_cdc_operation"` and `"_cdc_sequence_number"` columns created by Pipelines. A stream on the history table sees every appended row, including rows that represent source updates and deletes. Your job must interpret those operations, load the existing history, tolerate replay, and rebuild the current state after a source truncate or pipeline table reset.
### 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.
## Type mapping
Pipelines creates Snowflake columns with these mappings:
| Postgres type | Snowflake type |
| --------------------------------------- | ------------------------------- |
| `boolean` | `BOOLEAN` |
| `smallint`, `integer`, `bigint` | `SMALLINT`, `INTEGER`, `BIGINT` |
| `real`, `double precision` | `FLOAT`, `DOUBLE` |
| `date`, `time` | `DATE`, `TIME` |
| `timestamp`, `timestamp with time zone` | `TIMESTAMP_NTZ`, `TIMESTAMP_TZ` |
| `json`, `jsonb` | `VARIANT` |
| One-dimensional arrays | `ARRAY` |
| `oid` | `BIGINT` |
| Other types | `VARCHAR` |
Pipelines maps character and text types, `numeric`, `time with time zone`, `interval`, `uuid`, `bytea`, bit strings, and custom or unknown types to `VARCHAR`. It serializes these values instead of storing them as native Snowflake types. For `bytea`, the serialized value is a lowercase hexadecimal string.
Additional limits apply:
- Multi-dimensional arrays aren't supported. Non-default lower bounds on one-dimensional arrays aren't preserved.
- Non-finite `real` and `double precision` values are rejected. Non-finite `numeric` values are preserved as strings in `VARCHAR` columns.
- An uncompressed serialized row larger than 2 MiB is rejected.
- Source primary-key, unique, check, length, precision, and nullability constraints aren't copied. Only the two CDC metadata columns are `NOT NULL`.
## Schema change support
Pipelines supports:
- Adding, renaming, or dropping columns
- Adding or removing published columns on tracked tables
Replicated columns remain nullable in Snowflake, and changes to existing column defaults are not propagated. When a table is first created, Pipelines can copy compatible literal defaults. It can also copy string, numeric, or boolean literal defaults for columns added later in Postgres. Other defaults are omitted, although Postgres still supplies the source values through replication.
Schema changes also affect the stored history. Renaming a column changes its name in earlier events, dropping a column removes its historical values, and adding a column with a default can populate older rows.
Previously excluded columns are added without defaults, leaving historical events `NULL` for those columns. Removing a published column drops its destination values; adding it again does not restore them.
For type changes, unsupported changes, and interrupted schema changes, see the shared [schema-change behavior and recovery](/docs/guides/database/replication/pipelines#schema-change-support). Apart from [enabling change tracking](#materialize-with-a-dynamic-table), do not alter managed destination objects manually.
## Troubleshooting
| Symptom | What to check |
| ------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| Authentication fails | [Account identifier](#find-the-account-identifier), user, [key fingerprint and PEM](#generate-a-key-pair), and passphrase. Omit the passphrase for an unencrypted key. |
| Database or schema isn't found | Exact identifier case and `USAGE` permissions on both objects |
| Validation passes but table creation fails | [Grants and ownership](#prepare-snowflake-resources), including `CREATE TABLE` |
| Tables are created but writes fail | [SQL and streaming role alignment](#keep-the-sql-and-streaming-roles-aligned), table permissions, and network access to the account control and discovered Snowpipe ingest endpoints |
| Updates or deletes fail | [Source replica identity and published columns](#source-table-requirements) |
| A row is rejected | [Type and row-size limits](#type-mapping) and [reserved metadata columns](#append-only-change-history) |
| A schema change fails | [Supported changes](#schema-change-support); do not repair managed tables manually |
Use [pipeline monitoring](/docs/guides/database/replication/pipelines-monitoring) to inspect errors. For unresolved failures, [contact support](/dashboard/support/new) with the pipeline ID and error details.
## Additional resources
- [Snowflake key-pair authentication](https://docs.snowflake.com/en/user-guide/key-pair-auth)
- [Snowflake account identifiers](https://docs.snowflake.com/en/user-guide/admin-account-identifier)
- [Snowpipe Streaming default pipe](https://docs.snowflake.com/en/user-guide/snowpipe-streaming/snowpipe-streaming-pipe-object)
- [Snowpipe Streaming access control](https://docs.snowflake.com/en/user-guide/snowpipe-streaming/snowpipe-streaming-access-control)