Connect

snowflake_streaming

Ingest data from your pipeline into Snowflake using the Snowpipe Streaming classic architecture. This output does not work with the Snowpipe Streaming high-performance architecture.

To help you configure your own snowflake_streaming output, this page includes example data pipelines.

Introduced in version 4.39.0.

  • Common

  • Advanced

output:
  label: ""
  snowflake_streaming:
    account: "" # No default (required)
    user: "" # No default (required)
    role: "" # No default (required)
    database: "" # No default (required)
    schema: "" # No default (required)
    table: "" # No default (required)
    private_key: "" # No default (optional)
    private_key_file: "" # No default (optional)
    private_key_pass: "" # No default (optional)
    mapping: "" # No default (optional)
    init_statement: "" # No default (optional)
    schema_evolution:
      enabled: false # No default (required)
    batching:
      count: 0
      byte_size: 0
      period: ""
      check: ""
    max_in_flight: 4
output:
  label: ""
  snowflake_streaming:
    account: "" # No default (required)
    url: "" # No default (optional)
    user: "" # No default (required)
    role: "" # No default (required)
    database: "" # No default (required)
    schema: "" # No default (required)
    table: "" # No default (required)
    private_key: "" # No default (optional)
    private_key_file: "" # No default (optional)
    private_key_pass: "" # No default (optional)
    mapping: "" # No default (optional)
    init_statement: "" # No default (optional)
    schema_evolution:
      enabled: false # No default (required)
      ignore_nulls: true
      processors: [] # No default (optional)
    build_options:
      parallelism: 1
      chunk_size: 50000
    batching:
      count: 0
      byte_size: 0
      period: ""
      check: ""
      processors: [] # No default (optional)
    max_in_flight: 4
    channel_prefix: "" # No default (optional)
    channel_name: "" # No default (optional)
    offset_token: "" # No default (optional)
    commit_backoff:
      initial_interval: 32ms
      max_interval: 512ms
      max_elapsed_time: 60s
      multiplier: 2
    message_format: object
    timestamp_format: 2006-01-02T15:04:05.999999999Z07:00

Conversion of message data into Snowflake table rows

Message data conversion to Snowflake table rows is determined by the:

The following scenarios highlight how these three factors affect data written to the target table.

For reduced complexity, consider turning on schema evolution, which automatically creates and updates the Snowflake table schema based on message contents.

Scenario: Data and table schema match (schema evolution turned on or off)

An output message matches the existing table schema, and the schema_evolution.enabled field is set to true or false.

The target Snowflake table has two columns:

  • product_id (NUMBER)

  • product_code (STRING)

A pipeline generates the following message:

{"product_id": 521, "product_code": “EST-PR”}

In this scenario:

  • The JSON keys in the message ("product_id" and "product_code") match column names in the target Snowflake table.

  • The message values match the column data types. (If there was a data mismatch, the message would be rejected.)

  • Redpanda Connect inserts the message values into a new row in the target Snowflake table.

    product_id product_code

    521

    EST-PR

Scenario: Data and table schema mismatch (schema evolution turned on)

An output message includes schema updates, and the schema_evolution.enabled field is set to true.

The target Snowflake table has the same two columns as the previous scenario:

  • product_id (NUMBER)

  • product_code (STRING)

This time, the pipeline generates the following message:

{"product_batch": 11111, "product_color": “yellow”}

In this scenario:

  • The JSON keys ("product_batch" and "product_color") do not match column names in the target Snowflake table.

  • As schema evolution is enabled, Redpanda Connect adds two new columns to the target table with data types derived from the output message values. For more information about the mapping of data types, see Supported data formats for Snowflake columns.

  • Redpanda Connect inserts the message values into a new table row.

    product_id product_code product_batch product_color

    (null)

    (null)

    11111

    yellow

    You can configure processors to override the schema updates derived from the message values.

Scenario: Data and table schema mismatch (schema evolution turned off)

An output message includes schema updates, and the schema_evolution.enabled field is set to false.

The target Snowflake table has the same two columns:

  • product_id (NUMBER)

  • product_code (STRING)

The pipeline generates the same message as the previous scenario:

{"product_batch": 11111, "product_color": “yellow”}

In this scenario:

  • The JSON keys ("product_batch" and "product_color") do not match any existing column names.

  • Because schema evolution is turned off, Redpanda Connect ignores the extra column names and values and inserts a row of null values.

    product_id product_code

    (null)

    (null)

Supported data formats for Snowflake columns

The message data from your output must match the columns in the Snowflake table that you want to write data to. The following table shows you the column data types supported by Snowflake and how they correspond to the Bloblang data types in Redpanda Connect.

Snowflake column data type Bloblang data types

CHAR, VARCHAR

string

BINARY

string or bytes

NUMBER

number, or string where the string is parsed into a number

FLOAT, including special values, such as NaN (Not a Number), -inf (negative infinity), and inf (positive infinity)

number

BOOLEAN

bool, or number where a non-zero number is true

TIME, DATE, TIMESTAMP

timestamp, or number where the number is a converted to a Unix timestamp, or string where the string is parsed using RFC 3339 format

VARIANT, ARRAY, OBJECT

Any data type converted into JSON

GEOGRAPHY,GEOMETRY

Not supported

Authentication

You can authenticate with Snowflake using an RSA key pair. Either specify:

Performance

For improved performance, this output:

  • Sends multiple messages in parallel. You can tune the maximum number of in-flight messages (or message batches) with the field max_in_flight.

  • Sends messages as a batch. You can configure batches at both the input and output level. For more information, see Message Batching.

Batch sizes

Redpanda recommends that every message batch writes at least 16 MiB of compressed output to Snowflake. You can monitor batch sizes using the snowflake_compressed_output_size_bytes metric.

Metrics

This output emits the following metrics.

Metric name Description

snowflake_compressed_output_size_bytes

The size in bytes of each message batch uploaded to Snowflake.

snowflake_convert_latency_ns

The time taken to convert messages into the Snowflake column data types.

snowflake_serialize_latency_ns

The time taken to serialize the converted columnar data into a file for upload to Snowflake.

snowflake_build_output_latency_ns

The time taken to build the file that is uploaded to Snowflake. This metric is the sum of snowflake_convert_latency_ns + snowflake_serialize_latency_ns.

snowflake_upload_latency_ns

The time taken to upload the output file to Snowflake.

snowflake_register_latency_ns

The time taken to register the uploaded output file with Snowflake.

snowflake_commit_latency_ns

The time taken to commit the uploaded data updates to the target Snowflake table.

Fields

account

Use the format <orgname>-<account_name> where:

  • The <orgname> is the name of your Snowflake organization.

  • The <account_name> is the unique name of your account within your Snowflake organization.

To find the correct value for this field, run the following query in Snowflake:

WITH HOSTLIST AS
(SELECT * FROM TABLE(FLATTEN(INPUT => PARSE_JSON(SYSTEM$allowlist()))))
SELECT REPLACE(VALUE:host,'.snowflakecomputing.com','') AS ACCOUNT_IDENTIFIER
FROM HOSTLIST
WHERE VALUE:type = 'SNOWFLAKE_DEPLOYMENT_REGIONLESS';

Type: string

# Examples:
account: ORG-ACCOUNT

batching

Configure a batching policy.

Type: object

# Examples:
batching:
  byte_size: 5000
  count: 0
  period: 1s

# ---

batching:
  count: 10
  period: 1s

# ---

batching:
  check: this.contains("END BATCH")
  count: 0
  period: 1m

batching.byte_size

The maximum total size (in bytes) that a batch can reach before it is flushed. When the combined size of all messages in the batch reaches or exceeds this limit, the batch is immediately sent to the next stage (such as a processor or output).

Set to 0 to disable size-based batching. When disabled, messages are flushed based on other conditions (such as count or period).

Type: int

Default: 0

batching.check

A Bloblang query that returns a boolean value indicating whether a message should end a batch.

Type: string

Default: ""

# Examples:
check: this.type == "end_of_transaction"

batching.count

The number of messages at which the batch is flushed. Set to 0 to disable count-based batching.

Type: int

Default: 0

batching.period

The length of time after which an incomplete batch is flushed regardless of its size. This field accepts Go duration format strings such as 100ms, 1s, or 5s. Supported time units are ns, us, ms, s, m, and h.

Type: string

Default: ""

# Examples:
period: 1s

# ---

period: 1m

# ---

period: 500ms

batching.processors[]

A list of processors to apply to a batch as it is flushed. This allows you to aggregate and archive the batch however you see fit. All resulting messages are flushed as a single batch, so splitting the batch into smaller batches with these processors has no effect.

Type: array<processor>

# Examples:
processors:
  - archive:
      format: concatenate


# ---

processors:
  - archive:
      format: lines


# ---

processors:
  - archive:
      format: json_array

build_options

Options for optimizing the build of the output data that is sent to Snowflake. Monitor the snowflake_build_output_latency_ns metric to assess whether you need to update these options.

Requires version 4.40.0 or later.

Type: object

build_options.chunk_size

The number of table rows to submit in each chunk for processing.

Type: int

Default: 50000

build_options.parallelism

The maximum amount of parallel processing to use when building the output for Snowflake.

Type: int

Default: 1

channel_name

The channel name to use when connecting to a Snowflake table. Duplicate channel names cause errors and prevent multiple instances of Redpanda Connect from writing at the same time.

Redpanda Connect assumes that a message batch contains messages for a single channel, which means that interpolation is only executed on the first message in each batch. If your pipeline uses an input that is partitioned, such as an Apache Kafka topic, batch messages at the input level to make sure all messages in a batch are written to the same channel.

You can specify either the channel_name or channel_prefix, but not both. If neither field is populated, this output creates a channel name based on a table’s fully-qualified name, which results in a single stream per table.

Snowflake limits the number of streams per table to 10,000. If you need to use more than 10,000 streams, contact Snowflake support.

This field supports interpolation functions.

Requires version 4.45.0 or later.

Type: string

# Examples:
channel_name: partition-${!@kafka_partition}

channel_prefix

The prefix to use when creating a channel name for connecting to a Snowflake table. Adding a channel_prefix avoids the creation of duplicate channel names, which result in errors and prevent multiple instances of Redpanda Connect from writing at the same time.

You can specify either the channel_prefix or channel_name, but not both. If neither field is populated, this output creates a channel name based on a table’s fully-qualified name, which results in a single stream per table.

The maximum number of channels open at any time is determined by the value in the max_in_flight field.

Snowflake limits the number of streams per table to 10,000. If you need to use more than 10,000 streams, contact Snowflake support.

Type: string

# Examples:
channel_prefix: channel-${HOST}

commit_backoff

Control how frequently Snowflake is polled to check if data has been committed.

Requires version 4.82.0 or later.

Type: object

commit_backoff.initial_interval

The initial period to wait between status polls.

Type: string

Default: 32ms

commit_backoff.max_elapsed_time

The maximum total time to wait for data to be committed. If zero then no limit is used.

Type: string

Default: 60s

commit_backoff.max_interval

The maximum period to wait between status polls.

Type: string

Default: 512ms

commit_backoff.multiplier

The factor by which the poll interval grows on each attempt.

Type: float

Default: 2

database

The Snowflake database you want to write data to.

Type: string

# Examples:
database: MY_DATABASE

init_statement

Optional SQL statements to execute immediately after this output connects to Snowflake for the first time. This is a useful way to initialize tables before processing data.

Make sure your SQL statements are idempotent, so they do not cause issues when run multiple times after service restarts.

Type: string

# Examples:
init_statement: |

  CREATE TABLE IF NOT EXISTS mytable (amount NUMBER);

# ---

init_statement: |

  ALTER TABLE t1 ALTER COLUMN c1 DROP NOT NULL;
  ALTER TABLE t1 ADD COLUMN a2 NUMBER;

mapping

The Bloblang mapping to execute on each message.

Type: string

max_in_flight

The maximum number of messages to have in flight at a given time. For outputs that send messages in batches, this limit applies to message batches. Increase this value to improve throughput.

Type: int

Default: 4

message_format

The format to expect incoming messages from the rest of the pipeline.

Requires version 4.79.0 or later.

Type: string

Default: object

Option Summary

array

Messages are an array of values where the position in the array matches up the with ordinal of the column in snowflake

object

Messages are an object in JSON or bloblang where the key of the object is the column name in snowflake and the value is the value for the column

# Examples:
message_format: array

offset_token

The offset token to use for exactly-once delivery of data to a Snowflake table.

This output assumes that messages within a batch are in increasing order by offset token. When data is sent on a channel, the offset token of each message in the batch is compared to the latest token processed by the channel. If the offset token is less than the latest token, it’s assumed the message is a duplicate and is dropped. Messages must be delivered to the output in order, otherwise they are processed as duplicates and dropped.

Retried messages can also be seen as duplicates if later messages have succeeded in the meantime, so in most cases use a dead-letter queue to process failed messages. See the Ingesting data exactly once from Redpanda example.

If both offset tokens parse as base-10 integers (up to 64 bits), they are compared numerically, so a bare numeric token such as ${!@kafka_offset} does not need padding. Any other token, including numbers outside of the 64-bit integer range or a numeric value combined with a prefix or separator, falls back to a lexicographic comparison of its string representation. If you use one of these tokens, pad it so that it’s lexicographically ordered in its string representation.

For more information about offset tokens, see the Snowflake documentation.

This field supports interpolation functions.

Requires version 4.45.0 or later.

Type: string

# Examples:
offset_token: offset-${!"%016X".format(@kafka_offset)}

# ---

offset_token: postgres-${!@lsn}

private_key

The PEM-encoded private RSA key to use for authentication with Snowflake. You must specify a value for this field or the private_key_file field.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets.

Type: string

private_key_file

A .p8, PEM-encoded file to load the private RSA key from. You must specify a value for this field or the private_key field.

Type: string

private_key_pass

If the RSA key is encrypted, specify the RSA key passphrase.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets.

Type: string

role

The role of the user specified in the user field. The user’s role must have the required privileges to call the Snowpipe Streaming APIs. For more information about user roles, see the Snowflake documentation.

Type: string

# Examples:
role: ACCOUNTADMIN

schema

The schema of the Snowflake database you want to write data to.

Type: string

# Examples:
schema: PUBLIC

schema_evolution

Options to control schema updates when messages are written to the Snowflake table, such as adding columns as new fields appear in messages.

Type: object

schema_evolution.enabled

Whether schema evolution is enabled. When set to true, the Snowflake table is automatically created based on the schema of the first message written to it, if the table does not already exist. As new fields are added to subsequent messages in the pipeline, new columns are created in the Snowflake table. Any required columns are marked as nullable if new messages do not include data for them.

Type: bool

schema_evolution.ignore_nulls

When set to true and schema evolution is enabled, new columns that have null values are not added to the Snowflake table and do not trigger schema evolution. When set to false, null columns trigger schema migrations in Snowflake. Ignoring null values:

  • Prevents unnecessary schema changes caused by placeholder or incomplete data.

  • Avoids creating table columns with incorrect data types.

Redpanda does not recommend changing the default setting (true) unless you know the data type of null columns in advance.

Requires version 4.48.0 or later.

Type: bool

Default: true

schema_evolution.processors[]

A series of processors to execute when new columns are added to the Snowflake table. You can use these processors to:

  • Run side effects when the schema evolves.

  • Enrich the message with additional information to guide the schema changes.

For example, a processor could read the schema from the schema registry that a message was produced with and use that schema to determine the data type of the new column in Snowflake.

The input to these processors is an object with the value and name of the new column, the original message, and details of the Snowflake table the output writes to. The metadata remains the same as in the original message that triggered the schema update. For example: {"value": 42.3, "name":"new_data_field", "message": {"existing_data_field": 42, "new_data_field": "foo"}, "db": MY_DATABASE", "schema": "MY_SCHEMA", "table": "MY_TABLE"}

The output from the processors must be a single message that contains a string specifying the column data type to use, such as FLOAT, VARIANT, or NUMBER(38, 0). The output then runs an ALTER TABLE statement on the Snowflake table to add the column with the corresponding data type.

Requires version 4.46.0 or later.

Type: array<processor>

# Examples:
processors:
  - mapping: |-
      root = match this.value.type() {
        this == "string" => "STRING"
        this == "bytes" => "BINARY"
        this == "number" => "DOUBLE"
        this == "bool" => "BOOLEAN"
        this == "timestamp" => "TIMESTAMP"
        _ => "VARIANT"
      }

table

The Snowflake table you want to write data to.

This field supports interpolation functions.

Type: string

# Examples:
table: MY_TABLE

timestamp_format

The format to parse string values for TIMESTAMP, TIMESTAMP_LTZ and TIMESTAMP_NTZ columns. Should be a layout for time.Parse in Go.

Requires version 4.79.0 or later.

Type: string

Default: 2006-01-02T15:04:05.999999999Z07:00

url

Specify a custom URL to connect to Snowflake. This parameter overrides the default URL, which is generated from the value of the account field: https://<account>.snowflakecomputing.com.

Requires version 4.48.0 or later.

Type: string

# Examples:
url: https://org-account.privatelink.snowflakecomputing.com

user

Specify a user to run the Snowpipe Stream. To learn how to create a user, see the Snowflake documentation.

Type: string

Examples

Exactly once CDC into Snowflake

How to send data from a PostgreSQL table into Snowflake exactly once using Postgres Logical Replication.

If attempting to do exactly-once it’s important that rows are delivered in order to the output. Be sure to read the documentation for offset_token first. Removing the offset_token is a safer option that will instruct Redpanda Connect to use its default at-least-once delivery model instead.
input:
  postgres_cdc:
    dsn: postgres://foouser:foopass@localhost:5432/foodb
    schema: "public"
    slot_name: "my_repl_slot"
    tables: ["my_pg_table"]
    # We want very large batches - each batch will be sent to Snowflake individually
    # so to optimize query performance we want as big of files as we have memory for
    batching:
      count: 50000
      period: 45s
    # Prevent multiple batches from being in flight at once, so that we never send
    # a batch while another batch is being retried, this is important to ensure that
    # the Snowflake Snowpipe Streaming channel does not see older data - as it will
    # assume that the older data is already committed.
    checkpoint_limit: 1
output:
  snowflake_streaming:
    # We use the log sequence number in the WAL from Postgres to ensure we
    # only upload data exactly once, these are already lexicographically
    # ordered.
    offset_token: "${!@lsn}"
    # Since we're sending a single ordered log, we can only send one thing
    # at a time to ensure that we're properly incrementing our offset_token
    # and only using a single channel at a time.
    max_in_flight: 1
    account: "MYSNOW-ACCOUNT"
    user: MYUSER
    role: ACCOUNTADMIN
    database: "MYDATABASE"
    schema: "PUBLIC"
    table: "MY_PG_TABLE"
    private_key_file: "my/private/key.p8"

Ingesting data exactly once from Redpanda

How to ingest data from Redpanda with consumer groups, decode the schema using the schema registry, then write the corresponding data into Snowflake exactly once.

If attempting to do exactly-once it’s important that records are delivered in order to the output and correctly partitioned. Be sure to read the documentation for channel_name and offset_token first. Removing the offset_token is a safer option that will instruct Redpanda Connect to use its default at-least-once delivery model instead.
input:
  redpanda:
    topics: ["my_topic_going_to_snow"]
    consumer_group: "redpanda_connect_to_snowflake"
    # We want very large batches - each batch will be sent to Snowflake individually
    # so to optimize query performance we want as big of files as we have memory for
    fetch_max_bytes: 100MiB
    fetch_min_bytes: 50MiB
    partition_buffer_bytes: 100MiB
pipeline:
  processors:
    - schema_registry_decode:
        url: "redpanda.example.com:8081"
        basic_auth:
          enabled: true
          username: MY_USER_NAME
          password: "${TODO}"
output:
  fallback:
    - snowflake_streaming:
        # To ensure that we write an ordered stream each partition in kafka gets its own
        # channel.
        channel_name: "partition-${!@kafka_partition}"
        # Ensure that our offsets are lexicographically sorted in string form by padding with
        # leading zeros
        offset_token: offset-${!"%016X".format(@kafka_offset)}
        account: "MYSNOW-ACCOUNT"
        user: MYUSER
        role: ACCOUNTADMIN
        database: "MYDATABASE"
        schema: "PUBLIC"
        table: "MYTABLE"
        private_key_file: "my/private/key.p8"
        schema_evolution:
          enabled: true
    # In order to prevent delivery orders from messing with the order of delivered records
    # it's important that failures are immediately sent to a dead letter queue and not retried
    # to Snowflake. See the ordering documentation for the "redpanda" input for more details.
    - retry:
        output:
          redpanda:
            topic: "dead_letter_queue"

HTTP Server to push data to Snowflake

This example demonstrates how to create an HTTP server input that can receive HTTP PUT requests with JSON payloads, that are buffered locally then written to Snowflake in batches.

This example uses a buffer to respond to the HTTP request immediately, so it’s possible that failures to deliver data could result in data loss. See the documentation about buffers for more information, or remove the buffer entirely to respond to the HTTP request only once the data is written to Snowflake.
input:
  http_server:
    path: /snowflake
buffer:
  memory:
    # Max inflight data before applying backpressure
    limit: 524288000 # 50MiB
    # Batching policy, influences how large the generated files sent to Snowflake are
    batch_policy:
      enabled: true
      byte_size: 33554432 # 32MiB
      period: "10s"
output:
  snowflake_streaming:
    account: "MYSNOW-ACCOUNT"
    user: MYUSER
    role: ACCOUNTADMIN
    database: "MYDATABASE"
    schema: "PUBLIC"
    table: "MYTABLE"
    private_key_file: "my/private/key.p8"
    # By default there is only a single channel per output table allowed
    # if we want to have multiple Redpanda Connect streams writing data
    # then we need a unique channel prefix per stream. We'll use the host
    # name to get unique prefixes in this example.
    channel_prefix: "snowflake-channel-for-${HOST}"
    schema_evolution:
      enabled: true