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 VARIANT catch-all columns), row skips effectively do not occur — the ingestor rejects malformed envelope fields and VARIANT accepts anything. Skips arise only when you add typed playbook mappings into strictly-typed columns (for example, a string into a NUMBER column).
  • Column matching is case-insensitive. Add a VARIANT column (for example, properties, context, or traits) 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) plus properties, context, and traits as VARIANT. 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 404 error, which is visible.
  • Enabling ERROR_LOGGING = TRUE on the table is recommended. It is the only way skipped rows are recoverable: they can then be reviewed with ERROR_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 = TRUE on the table auto-adds unmatched keys as new type-inferred columns (requires the EVOLVE SCHEMA privilege), and a custom pipe (via PIPE) can apply server-side casts or transforms.
  • Required Snowflake grants are USAGE on the database and schema, INSERT on the table, and OPERATE on 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 USAGE on the database and schema, INSERT on the destination table, and OPERATE on the pipe.
  • Create or identify the destination database, schema, and table. The table must exist before deployment; enabling ERROR_LOGGING = TRUE on 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 ParameterDescription
ACCOUNTSnowflake 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.
USERSnowflake user that authenticates and loads events.
PRIVATE_KEYPrivate 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.
DATABASESnowflake database containing the destination table.
SCHEMASchema containing the destination table.
TABLEDestination table that receives events.
PIPEEnter <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 VARIANT catch-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 KeyDescriptionExpected Input
Expression OutputObject: 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_DELIVERED before 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_HISTORY account usage view.
  • If ERROR_LOGGING = TRUE is enabled on the table, review skipped rows with the ERROR_TABLE() table function, filtering on error_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 VARIANT catch-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