Snowflake Streaming
What Is Snowflake Streaming?
Snowflake is a cloud data platform that lets organizations store, process, and analyze large volumes of data across multiple clouds. Snowpipe Streaming is Snowflake's low-latency ingestion service, enabling data to be loaded into Snowflake tables as a continuous stream rather than in scheduled file batches. Its high-performance architecture ingests rows through a REST API using durable channels and offset-based delivery, making it well suited for event-driven data pipelines. This allows teams to centralize behavioral and transactional event data in Snowflake for analytics, reporting, and downstream activation.
Product Type: Data Warehouse
Integration Type: Starter Kit
Event Source Type: Web and/or Mobile Browser and/or Mobile App
Event Scope: Full-Funnel Events
Capabilities
- Streams event data directly into a Snowflake table using Snowpipe Streaming.
- Customizable delivery — events land as-is by default (matched to table columns by name), and you can add playbook transforms to shape events to your table's schema.
- Authenticates using key-pair JWT, exchanged for a short-lived scoped token.
- Automatically chunks oversized batches into smaller appends on the same channel.
- Preserves event ordering within a channel.
- Resumes safely on retry using offset tokens, avoiding duplicates in the normal path.
- Operates with a minimal set of Snowflake privileges.
Considerations
- The destination table schema is the contract for this integration. Because events are sent as-is, any top-level field without a matching column is dropped, and any value that cannot cast to its column's type is skipped. Neither condition fails the batch or appears as an error.
- In the common passthrough setup (events as-is plus
VARIANTcatch-all columns), row skips effectively do not occur — the ingestor rejects malformed envelope fields andVARIANTaccepts anything. Skips arise only when you add typed playbook mappings into strictly-typed columns (for example, a string into aNUMBERcolumn). - Column matching is case-insensitive. Add a
VARIANTcolumn (for example,properties,context, ortraits) as a catch-all so unmatched and nested fields are retained. - Recommended table shape for Analytics.js events: typed columns for the envelope fields (
type,event,anonymousId,userId,messageId,timestamp) plusproperties,context, andtraitsasVARIANT. Widen or add typed columns deliberately. - The destination table must exist before deployment. The pipe is auto-created, but the table is not. A missing or unauthorized table returns a non-retryable
404error, which is visible. - Enabling
ERROR_LOGGING = TRUEon the table is recommended. It is the only way skipped rows are recoverable: they can then be reviewed withERROR_TABLE(), which returns the full value, the affected column, and the original event payload. Without it, skipped rows cannot be recovered, and the channel status reason redacts the value and type. - Row-level skips surface only as a cumulative per-channel counter on later batches, in forwarder logs, or via
ERROR_TABLE()— not on the event that caused the skip. - Two Snowflake-side options reduce silent drops:
ENABLE_SCHEMA_EVOLUTION = TRUEon the table auto-adds unmatched keys as new type-inferred columns (requires theEVOLVE SCHEMAprivilege), and a custom pipe (viaPIPE) can apply server-side casts or transforms. - Required Snowflake grants are
USAGEon the database and schema,INSERTon the table, andOPERATEon the pipe. No warehouse or storage stage is required. - Authentication uses an RSA key-pair JWT. You generate the key pair, assign the public key to your Snowflake user (
ALTER USER <user> SET RSA_PUBLIC_KEY=…), and paste the matching private key into the integration's Private Key field. MetaRouter stores that private key securely and uses it only to sign the JWT — MetaRouter does not generate or hold any credentials of its own. - Per-append size limits are 16 MB total and 4 MB for rows. Oversized batches are automatically chunked into smaller appends on the same channel. Events are capped at 250 KB at ingest.
- Per-channel sustained throughput is approximately 20 MB/s and 10 requests per second.
- Per-pipe limits are 2,000 active channels (raisable via Snowflake support), 10 pipes per table, 10 GB/s aggregate per table, and 10,000 pipes per account.
- Channels auto-delete after 30 days of inactivity and reopen transparently on the next batch.
- Events are delivered in batches (default 500); a failed row is skipped individually rather than rejecting the whole batch. Batch size is a platform-level setting MetaRouter can tune on request — not a self-serve field in the UI.
- Delivery is idempotent: offset tokens allow a retried batch to resume from the committed point, so retries do not create duplicates in the normal path.
- This integration uses the Snowpipe Streaming high-performance architecture, generally available in AWS, Azure, and GCP commercial regions (not government regions). The account identifier may be in locator or organization format.
Limitations
- A delivered status means Snowflake accepted the rows for processing, not that they are committed or queryable. Commit is asynchronous (approximately 5–10 seconds documented, and around 18–20 seconds observed in testing), and latency can rise briefly during throughput spikes.
- Column-level and row-level drops are not surfaced as batch errors and cannot be traced to the specific event that caused them from the event debugger alone.
- In a rare fast-retry edge case, buffered-but-not-yet-committed rows from an earlier acknowledged batch can be discarded if a retry reopens the channel before that batch commits.
- Very large single appends can reach the forwarder HTTP client timeout (60 seconds). This is only relevant for very large events on constrained egress.
- There is no cross-channel ordering guarantee. Very high-volume single write keys are bounded by the per-channel throughput ceilings.
- This integration uses the high-performance architecture and is separate from the file-based Snowflake batch integration; it does not replace it.
Starter Kit Setup Guide
1. Gather Credentials
- Generate an RSA key pair and assign the public key to the Snowflake user MetaRouter will authenticate as. You'll paste the matching private key into the integration in a later step. (See the SQL block below.)
- Grant the user
USAGEon the database and schema,INSERTon the destination table, andOPERATEon the pipe. - Create or identify the destination database, schema, and table. The table must exist before deployment; enabling
ERROR_LOGGING = TRUEon the table is recommended. - Identify the pipe that loads events into the table. The default streaming pipe is named
<TABLE>-STREAMING; a custom pipe can be used to apply server-side casts or transforms. - Collect the account identifier, user, database, schema, table, and pipe name. The organization account format (for example,
myorg-account123) is recommended; the locator format is also accepted.
Use the following SQL to set up the Snowflake side (replace the <...> placeholders):
-- 1. Generate an RSA key pair locally (NOT in Snowflake):
-- openssl genrsa 2048 | openssl pkcs8 -topk8 -inform PEM -out rsa_key.p8 -nocrypt
-- openssl rsa -in rsa_key.p8 -pubout -out rsa_key.pub
-- Paste the full rsa_key.p8 PEM (with BEGIN/END lines) into the integration's Private Key field.
-- 2. Assign the public key to the Snowflake user
-- (paste the rsa_key.pub body only — strip the BEGIN/END PUBLIC KEY lines):
ALTER USER <user> SET RSA_PUBLIC_KEY='MIIBIjANBgkq...';
-- 3. Create the destination table (recommended Analytics.js shape).
-- Column names match AJS keys case-insensitively; VARIANT columns catch nested objects.
CREATE TABLE <db>.<schema>.<table> (
type STRING,
event STRING,
anonymousId STRING,
userId STRING,
messageId STRING,
timestamp TIMESTAMP_NTZ,
properties VARIANT,
context VARIANT,
traits VARIANT
);
-- 4. Enable per-row error recovery (strongly recommended, set at setup time):
ALTER TABLE <db>.<schema>.<table> SET ERROR_LOGGING = TRUE;
-- 5. Grant the role MetaRouter's user authenticates with:
GRANT USAGE ON DATABASE <db> TO ROLE <role>;
GRANT USAGE ON SCHEMA <db>.<schema> TO ROLE <role>;
GRANT INSERT ON TABLE <db>.<schema>.<table> TO ROLE <role>;
GRANT OPERATE ON PIPE <db>.<schema>."<TABLE>-STREAMING" TO ROLE <role>;
-- Optional: auto-add unmatched fields as new columns instead of dropping them:
ALTER TABLE <db>.<schema>.<table> SET ENABLE_SCHEMA_EVOLUTION = TRUE;
GRANT EVOLVE SCHEMA ON TABLE <db>.<schema>.<table> TO ROLE <role>;
-- Optional: custom pipe for server-side transforms (then set PIPE to this name):
CREATE PIPE <db>.<schema>.<pipe_name> AS
COPY INTO <db>.<schema>.<table>
FROM (SELECT <transforms> FROM TABLE(DATA_SOURCE(TYPE => 'STREAMING')));2. Add a Snowflake Streaming Integration
- From the integration library, add a Snowflake Streaming integration. Then, fill out the Connection Parameters:
| Connection Parameter | Description |
|---|---|
ACCOUNT | Snowflake account identifier: locator form (for example, xy12345.us-east-1 or xy12345.eu-central-2.aws, including region and cloud) or organization form MYORG-MYACCOUNT. |
USER | Snowflake user that authenticates and loads events. |
PRIVATE_KEY | Private key (PEM) used for key-pair authentication. You generate it and paste it here; MetaRouter stores it securely and uses it only to sign the JWT. |
DATABASE | Snowflake database containing the destination table. |
SCHEMA | Schema containing the destination table. |
TABLE | Destination table that receives events. |
PIPE | Enter <TABLE>-STREAMING for the Snowflake-managed default pipe (auto-created on first delivery), or the name of a custom pipe created for server-side transforms. |
3. Configure Event Mapping
- The kit delivers events as-is by default, matching JSON fields to your table columns by name (case-insensitive). You own the mapping: add playbook transforms to reshape events, and design your table (typed columns plus
VARIANTcatch-alls) to decide what is stored. See the recommended table shape under Considerations.
4. Deploy to Pipeline
- In the Pipelines tab, add your Snowflake Streaming integration.
- Select the correct integration revision.
- Click Add Integration to finalize deployment.
Event Mappings
Global
Global mappings will be applied to all events. If your parameter names do not match the Expected Inputs provided, you will need to overwrite the Inputs provided with your own.
| Output Key | Description | Expected Input |
|---|---|---|
| Expression Output | Object: The complete event is forwarded to the destination table without field-level transformation. | Expression – returns the full input event |
Event Specific
This integration does not use per-event mappings. All events are forwarded to Snowflake exactly as received, and the destination table schema determines which fields are stored.
Required & Recommended Identifiers
This integration forwards events to a destination table as-is and does not perform user matching or attribution. No specific identifiers are required for matching. The destination table schema determines which fields are stored; see Considerations for the recommended table shape.
Integration Validation
- Confirm events are landing by querying the destination table in Snowflake (for example, in a Snowsight worksheet) and checking for recently ingested rows. Allow roughly 20–30 seconds after an event shows
STATUS_DELIVEREDbefore querying — commit is asynchronous, so rows are not queryable the instant delivery succeeds. - Monitor ingestion progress and error counts for the pipe's channels using Snowflake's channel status, or review historical trends and error patterns in the
SNOWPIPE_STREAMING_CHANNEL_HISTORYaccount usage view. - If
ERROR_LOGGING = TRUEis enabled on the table, review skipped rows with theERROR_TABLE()table function, filtering onerror_metadata:service = 'snowpipe_streaming'. Each record includes the error message, the affected column, and the original payload for diagnosis and reprocessing. - Verify that expected columns are populated to confirm no fields are being silently dropped. If fields are missing, add a
VARIANTcatch-all column or enable schema evolution on the table.
After go-live, if you suspect rows were skipped, review them with the error table. This requires ERROR_LOGGING = TRUE to have been set at setup time (see the setup SQL block).
-- Review rows Snowflake skipped (unredacted value, affected column, full original payload):
SELECT * FROM ERROR_TABLE(<db>.<schema>.<table>);ERROR_TABLE() is a table function, so it does not appear in the Snowsight object browser. Without ERROR_LOGGING = TRUE, skipped rows are unrecoverable and the channel-status reason redacts the value and type.
Additional Snowflake Streaming Documentation
- Snowpipe Streaming high-performance overview
- REST API endpoints
- REST cURL/JWT tutorial
- Channels and offset tokens
- PIPE object
- Limitations
- Error logging / ERROR_TABLE
- Cost model
Updated 38 minutes ago