What you will be able to do
- Explain where the Snowflake Connector for Kafka runs and which versions exist
- Identify the core connection properties needed to configure the connector
- Predict the Snowflake table name the connector derives from a Kafka topic name
- Contrast the classic (v3) two-VARIANT-column schema with v4 default-pipe and user-defined-pipe modes
- Control RECORD_METADATA content and its per-record overhead
Key concept
Connector placement — Each Snowflake connector runs in a specific place: the Kafka connector inside a Kafka Connect cluster, the Spark connector inside a Spark cluster, and the Python connector inside your own application. A native connector is different: it is a connector application built with the Snowflake Native App Framework, so it is installed and runs inside your Snowflake account and pulls data in from the source system. Where a connector runs tells you what you have to install and operate yourself.
1.Where the Kafka connector runs
Apache Kafka uses a publish/subscribe model. Producers write messages to topics, and consumers read them asynchronously, so a subscriber never has to be connected directly to the publisher. To scale, topics are split into partitions. Kafka Connect is the framework that links Kafka to external systems such as databases. A Kafka Connect cluster is separate from the Kafka cluster itself, and it is where connectors run and scale out.
The Snowflake Connector for Kafka is one of those connectors. It reads one or more topics and loads the records into Snowflake tables. The typical pattern is one topic feeding one table, with each Kafka message becoming one row. Snowflake ships a build for the Confluent package and a build for open-source Apache Kafka. Confluent Cloud also offers a hosted version, so in that case you don't run the Connect workers yourself.
The classic connector loads data with one of two methods: Snowpipe or Snowpipe Streaming. You also need to know which generation of the connector you are using. The classic connector (v3 and earlier) is still fully supported, but Snowflake's documentation says it is planned for future deprecation. A formal announcement was planned for mid-2026, followed by an 18-month migration window. Snowflake recommends v4 for all new implementations. The two generations handle tables differently, so the rest of this page labels each behaviour with its version.
Classic connector buffering. The classic connector buffers records in memory per Kafka partition and flushes them to an internal stage as data files when a threshold is reached. The classic install documentation lists three thresholds: buffer.count.records (default 10000 records), buffer.flush.time (default 120 seconds between flushes) and buffer.size.bytes (default 5000000, or 5 MB). Records are compressed when they are written to data files, so the size of the records in the buffer can be larger than the size of the files created from them.
| Aspect | Classic connector (v3 and earlier) | Snowflake Connector for Kafka (v4) |
|---|---|---|
| Status | Fully supported; planned for future deprecation | Recommended for all new implementations |
| Load method | Snowpipe or Snowpipe Streaming | Snowpipe Streaming through a pipe (default {tableName}-STREAMING or user-defined) |
| Default table shape | Two VARIANT columns: RECORD_CONTENT and RECORD_METADATA | Columns matching the record's first-level JSON keys |
Checkpoint 1 of 6· Check yourself
An architect is sizing the infrastructure for the Snowflake Kafka connector, which they will self-manage. Where do the connector processes run?
The connector is a Kafka Connect component. Kafka Connect runs as its own cluster, apart from the Kafka brokers.
“A Kafka Connect cluster is a separate cluster from the Kafka cluster.”Source: docs.snowflake.com
2.Configuring the connection
Connector configuration is a set of properties. Some tell the connector where to write in Snowflake. Others tell Kafka Connect how to deserialize records. The Snowflake-side properties identify the account, the user, the role used for inserts, and the database and schema that hold the target tables. The URL must include your account identifier, and the https:// protocol and port number are optional.
The converter properties belong to Kafka Connect. For JSON records, the value converter is the community JsonConverter. For Avro records with a Schema Registry, it is Confluent's AvroConverter. One restriction to remember: when snowflake.enable.schematization is true (the default), StringConverter and ByteArrayConverter aren't supported as value converters.
| Property | What it sets |
|---|---|
| snowflake.url.name | URL for your Snowflake account; must include the account identifier |
| snowflake.user.name | User login name for the Snowflake account |
| snowflake.role.name | Role the connector uses to insert data into the table |
| snowflake.database.name / snowflake.schema.name | Database and schema containing the target table |
| value.converter | Record format, e.g. org.apache.kafka.connect.json.JsonConverter or io.confluent.connect.avro.AvroConverter |
| key.converter | Kafka record key converter; required by the Kafka Connect Platform |
Checkpoint 2 of 6· Match them up
Match each connector property to what it controls
Tap a term, then the definition that fits it.
The snowflake.* properties point the connector at Snowflake objects. The converters tell Kafka Connect how to read each record.
“The name of the role that the connector will use to insert data into the table.”Source: docs.snowflake.com
Sources3
3.From topic name to table name
After the connector is connected, it has to decide which table each topic writes to. There are two modes. In static mapping, the table name is derived from the topic name alone. In explicit mapping, you list topic-to-table pairs in snowflake.topic2table.map. Without that parameter, the connector always derives the name from the topic.
The classic install documentation describes the map format. Each topic and its table are separated by a colon, and pairs are comma-separated. Several topics can point at one table, for example topic1:low_range,topic2:low_range. Regular expressions are allowed to define topics, as with topics.regex, so topic[0-4]:low_range,topic[5-9]:high_range is valid. The expressions cannot be ambiguous: any matched topic must match only a single target table.
The derivation rules are mechanical. A topic name that is a valid Snowflake identifier becomes the table name in uppercase. Invalid characters are replaced with underscores, and a name that starts with something other than a letter or underscore gets an underscore prepended. Replacing characters can make two different topics produce the same name, so whenever the connector adjusts a name it appends an underscore and a hash code. The v4 docs give the example that topic my-topic.data becomes MY_TOPIC_DATA_<hash>. Snowflake's advice is to choose topic names that are already valid identifiers so table names stay predictable.
Checkpoint 3 of 6· Check yourself
A connector reads topics numbers+x and numbers-x with no snowflake.topic2table.map configured. What happens to the table names?
After invalid characters are replaced, both names become NUMBERS_X. To avoid accidental duplication, the connector adds a hash suffix.
“To avoid accidental duplication of table names, the connector appends a suffix to the table name.”Source: docs.snowflake.com
4.What lands in the table: classic columns versus v4 pipes
Classic connector. By default, every table it loads has two VARIANT columns. RECORD_CONTENT holds the Kafka message (JSON or Avro), and RECORD_METADATA holds details such as the topic, partition and offset. The message is not parsed or split into columns. If Snowflake creates the table, it has only those two columns. If you create the table yourself, you can add more columns, but they must allow NULL because the connector never supplies values for them. You read the data with VARIANT path syntax:
select
record_metadata:CreateTime,
record_content:ID
from table1
where record_metadata:topic = 'PressureOverloadWarning';v4 connector. You don't have to create a table or pipe in advance. The connector creates a table whose columns match the JSON keys and uses a default pipe named {tableName}-STREAMING. That pipe maps first-level keys to columns by name, ignoring case. Nested objects and arrays should go into VARIANT columns. Keys that don't match a column are ignored unless ENABLE_SCHEMA_EVOLUTION is enabled on the table, in which case the connector adds the new columns.
If you need renamed columns, casts, masking or filtering, use user-defined pipe mode. You create a pipe with the same name as the destination table and put the transformation in its COPY INTO, reading from the streaming data source. When choosing a pipe, the connector looks for a pipe with the table's name and falls back to the default if there isn't one. User-defined pipes are not available with client-side validation (snowflake.validation=client_side).
Checkpoint 4 of 6· Put it in order
Put the v4 connector's steps for choosing a destination and pipe in order
- 1.Check whether a pipe exists with the same name as the destination table
- 2.If no such user-created pipe exists, use the default pipe named {tableName}-STREAMING
- 3.Derive the destination table name from the topic name (or the topic2table map)
Pipe lookup depends on the table name, so the table name comes first. A user-created pipe with that name wins, and the default pipe is the fallback.
“The connector checks if a pipe exists with the same name as the destination table name.”Source: docs.snowflake.com
Checkpoint 5 of 6· Fill the gap
Which data source type completes this user-defined pipe for the v4 connector?
CREATE PIPE ORDERS AS
COPY INTO ORDERS
FROM (
SELECT
$1:order_id::STRING,
$1:customer_name,
$1:order_total::STRING,
$1:isPaid::STRING
FROM TABLE(DATA_SOURCE(TYPE => ' ? '))
);The v4 connector delivers records through Snowpipe Streaming, so a user-defined pipe reads from DATA_SOURCE(TYPE => 'STREAMING').
Source: docs.snowflake.comSources4
5.Controlling RECORD_METADATA
In both generations, the connector records where each row came from. RECORD_METADATA carries the topic, the Kafka partition (not a Snowflake micro-partition), the offset, and CreateTime or LogAppendTime in milliseconds since the epoch. It can also carry the message key and headers. With Snowpipe Streaming it includes SnowflakeConnectorPushTime, the moment the record entered the ingestion buffer, which is useful for measuring end-to-end latency. In a v4 user-defined pipe, you can pull individual fields out with $1:RECORD_METADATA.
The key is only stored if key.converter is set to org.apache.kafka.connect.storage.StringConverter. With any other converter, the connector ignores keys. Metadata also has a cost: about 150 bytes per record, depending on the length of the topic name. That can be a large share of a small payload and it affects billing. Boolean flags control which fields are included.
| Property | Includes |
|---|---|
| snowflake.metadata.topic | Topic name |
| snowflake.metadata.offset.and.partition | Offset and partition |
| snowflake.metadata.createtime | Kafka record timestamp |
| snowflake.metadata.all | All available metadata |
Checkpoint 6 of 6· Check yourself
A topic carries very small IoT payloads, and ingestion cost is higher than expected. Nobody downstream uses RECORD_METADATA. What change does the documentation suggest?
RECORD_METADATA adds roughly 150 bytes per record, so turning it off with snowflake.metadata.all=false removes that overhead.
“If you don’t need metadata, you can reduce this overhead by setting snowflake.metadata.all=false.”Source: docs.snowflake.com
Sources4
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.The classic Kafka connector parses each JSON message into separate typed columns.Why is that wrong?
By default, the classic connector stores the whole message in a single VARIANT column (RECORD_CONTENT) without parsing it. Column-per-key mapping is v4 default-pipe behaviour.
Covered in What lands in the table: classic columns versus v4 pipes
2.Message keys always appear in RECORD_METADATA, whatever key.converter is set to.Why is that wrong?
Keys are stored only when key.converter is set to org.apache.kafka.connect.storage.StringConverter. Otherwise they are ignored.
Covered in Controlling RECORD_METADATA
3.Columns you add to a classic connector table can be NOT NULL, because the connector fills them.Why is that wrong?
The connector only writes its own columns. Any extra columns must allow NULL.
Covered in What lands in the table: classic columns versus v4 pipes
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Recommendation: Use the Snowflake Connector for Kafka (v4) for all new implementations.”
↩︎ Where the Kafka connector runs“With Snowflake, the typical pattern is that one topic supplies messages (rows) for one Snowflake table.”
↩︎ Where the Kafka connector runs“Snowflake plans to issue a formal deprecation announcement in mid-2026.”
↩︎ Where the Kafka connector runs“The Kafka connector is designed to run in a Kafka Connect cluster”
↩︎ Key concept“The data is not parsed, and the data is not split into multiple columns in the Snowflake table.”
↩︎ Exam trap 1“any additional columns must allow NULL values because data from the connector does not include values for those columns”
↩︎ Exam trap 3“The current version of the Kafka connector is limited to loading data into Snowflake.”
↩︎ Prediction“A Kafka Connect cluster is a separate cluster from the Kafka cluster.”
↩︎ Checkpoint“To avoid accidental duplication of table names, the connector appends a suffix to the table name.”
↩︎ Checkpoint - 2.
“Number of records buffered in memory per Kafka partition before ingesting to Snowflake. The default value is 10000 records.”
↩︎ Where the Kafka connector runs“where the flush is from the Kafka’s memory cache to the internal stage. The default value is 120 seconds.”
↩︎ Where the Kafka connector runs“The topic configuration allows use of regular expressions to define topics, just as the use of topics.regex does.”
↩︎ From topic name to table name“any matched topic must match only a single target table”
↩︎ From topic name to table name - 3.
“The URL for accessing your Snowflake account. This URL must include your account identifier.”
↩︎ Configuring the connection“When snowflake.enable.schematization=true (the default), StringConverter and ByteArrayConverter aren’t supported as value converters.”
↩︎ Configuring the connection“The name of the role that the connector will use to insert data into the table.”
↩︎ Checkpoint - 4.
“If you do not configure the snowflake.topic2table.map parameter, the connector always derives the table names from the topic name.”
↩︎ From topic name to table name“For example, the topic my-topic.data becomes MY_TOPIC_DATA_<hash>”
↩︎ From topic name to table name“By default you don’t have to create a table or pipe before ingestion begins.”
↩︎ What lands in the table: classic columns versus v4 pipes“If keys from the JSON don’t match the table columns, the connector ignores the keys unless ENABLE_SCHEMA_EVOLUTION is enabled on the table”
↩︎ What lands in the table: classic columns versus v4 pipes“When using client-side validation (snowflake.validation=client_side), only the default pipe mode is supported.”
↩︎ What lands in the table: classic columns versus v4 pipes“RECORD_METADATA adds approximately 150 bytes of overhead per record, depending on the topic name length.”
↩︎ Controlling RECORD_METADATA“the key.converter parameter in the connector configuration properties must be set to org.apache.kafka.connect.storage.StringConverter”
↩︎ Exam trap 2“The connector checks if a pipe exists with the same name as the destination table name.”
↩︎ Checkpoint“If you don’t need metadata, you can reduce this overhead by setting snowflake.metadata.all=false.”
↩︎ Checkpoint