概述
Snowflake Sink Connector 将 Kafka Topic 中的记录持续写入 Snowflake 表,位于 Kafka 与 Snowflake 之间的数据写入链路中。它使用 Snowpipe Streaming 将每个 Topic 分区的数据写入对应的 Snowflake 摄取通道,可根据 Topic 名称自动确定目标表,也可通过映射将 Topic 路由到指定表。 Connector 可以把结构化记录展开为 Snowflake 列,也可以将记录写入RECORD_CONTENT 和 RECORD_METADATA 列,适合持续写入业务事件、日志和变更数据。
前置条件
- Snowflake 中的目标数据库、Schema、用户和角色已创建,并具备访问目标对象、创建或写入表以及使用 Snowpipe Streaming 所需的权限。
- JWT 认证需要 RSA 私钥;OAuth 认证需要客户端 ID、客户端密钥以及所选流程所需的令牌或端点。
- 使用已有表时,表结构应与记录转换方式和元数据配置相容;使用托管 Iceberg 表时,还需准备外部卷和目录配置。
授权许可
使用 Apache License 2.0。快速开始
请提前准备 Connect Cluster、Kafka Topic 和 Snowflake 目标数据库与 Schema,确认 Connect Worker 能访问 Kafka 和 Snowflake,并具备相应权限。Connector 的创建和管理方式请参考 AutoMQ 的管理 Connector。orders Topic,在目标 Schema 中自动创建缺失的 Snowflake 表,并将记录写入 RECORD_CONTENT 和 RECORD_METADATA 列。
配置
连接与认证
snowflake.url.name
Snowflake 账号 URL。该配置没有可用的运行时默认值。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:是
snowflake.user.name
用于建立 JWT 或 OAuth Snowflake 会话的用户。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:是
snowflake.private.key
默认使用 snowflake_jwt 认证时必填;可以通过 Kafka Config Provider 引用提供。
- 类型:
PASSWORD - 默认值:空字符串
- 重要级别:高
- 必填:条件必填
snowflake.private.key.passphrase
仅在 RSA 私钥已加密时设置;可以使用 Config Provider 引用。
- 类型:
PASSWORD - 默认值:空字符串
- 重要级别:低
- 必填:否
snowflake.database.name
目标 Snowflake 数据库必须已存在且当前角色可访问。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:是
snowflake.schema.name
目标 Schema 必须已存在且当前角色可访问。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:是
snowflake.role.name
指定 Snowflake 会话角色。启用 OAuth scope 但未显式配置 scope 时,还用于生成 session:role:<role>。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:是
snowflake.authenticator
选择 JWT 或 OAuth 认证,两种方式所需的凭据互斥。有效值为 snowflake_jwt、oauth。
- 类型:
STRING - 默认值:snowflake_jwt
- 重要级别:低
- 必填:否
snowflake.oauth.client.id
使用 OAuth 时必填;通过 Config Provider 引用时,实时凭据校验会延后。
- 类型:
STRING - 默认值:空字符串
- 重要级别:高
- 必填:条件必填
snowflake.oauth.client.secret
使用 OAuth 时必填;通过 Config Provider 引用时,实时凭据校验会延后。
- 类型:
PASSWORD - 默认值:空字符串
- 重要级别:高
- 必填:条件必填
snowflake.oauth.refresh.token
空值使用 OAuth client_credentials 流程;非空值使用 refresh token 流程。
- 类型:
PASSWORD - 默认值:空字符串
- 重要级别:高
- 必填:否
snowflake.oauth.token.endpoint
未设置时根据 Snowflake 账号 URL 推导 token endpoint。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 必填:否
snowflake.oauth.include.scope
false 不发送 scope;true 发送显式 scope,未设置时发送由角色生成的 session:role:<role>。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
snowflake.oauth.scope
仅在 snowflake.oauth.include.scope=true 时生效;此时留空将使用由角色生成的 scope。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
网络代理
jvm.proxy.host
非空值会设置进程级 HTTP 和 HTTPS JVM 代理属性;必须与 jvm.proxy.port 同时配置。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:条件必填
jvm.proxy.port
作为 JVM 代理端口字符串传入;必须与 jvm.proxy.host 同时配置。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:条件必填
jvm.nonProxy.hosts
代理主机和端口启用时,以 | 连接到现有的 http.nonProxyHosts JVM 属性。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
jvm.proxy.username
仅在代理主机和端口已启用时使用;必须与 jvm.proxy.password 同时配置。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:条件必填
jvm.proxy.password
仅在代理主机和端口已启用时使用;必须与 jvm.proxy.username 同时配置。
- 类型:
PASSWORD - 默认值:空字符串
- 重要级别:低
- 必填:条件必填
表与数据模型
snowflake.metadata.all
元数据总开关。设为 false 时忽略各字段开关并丢弃记录元数据;托管 Iceberg 表必须保持为 true。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.metadata.createtime
控制创建时间元数据;托管 Iceberg 表必须保持为 true。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.metadata.topic
控制 Topic 元数据;托管 Iceberg 表必须保持为 true。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.metadata.offset.and.partition
同时控制 Kafka offset 和分区元数据;托管 Iceberg 表必须保持为 true。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.streaming.metadata.connectorPushTime
控制 Connector 推送时间元数据;托管 Iceberg 表必须保持为 true。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.feature.structured.headers
true 保留转换后的结构化 Header 类型;false 将 Header 值扁平化为字符串。启用后可能与依赖旧元数据表示的下游不兼容。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
兼容性与迁移
snowflake.streaming.validate.compatibility.with.classic
默认 true 是 v3 迁移保护门禁。保持启用时,必须使用 snowflake.validation=client_side,将两个 snowflake.compatibility.* 规范化开关设为 true,并显式设置 snowflake.enable.schematization 和 snowflake.streaming.classic.offset.migration;迁移模式为 strict 或 best_effort 时,还必须显式设置 snowflake.streaming.classic.offset.migration.include.connector.name。新建 v4 Connector 且不需要 v3 兼容时,可显式设为 false。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:高
- 必填:否
Topic 路由、校验与迁移
snowflake.topic2table.map
以逗号分隔 topic:table 映射。带引号的表名保留大小写,未加引号的表名转为大写;空值按 Topic 名生成表名。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
snowflake.validation
有效值为 server_side、client_side。server_side 要求目标表启用错误日志;托管 Iceberg 表不支持 client_side。
- 类型:
STRING - 默认值:server_side
- 重要级别:高
- 必填:否
snowflake.streaming.classic.offset.migration
有效值为 skip、best_effort、strict。仅当 SSv2 通道没有已提交 offset 时才尝试读取 v3 Classic 通道:strict 在旧通道不存在时失败,best_effort 回退到 Kafka consumer group offset,skip 不读取旧通道。
- 类型:
STRING - 默认值:skip
- 重要级别:高
- 必填:条件必填
snowflake.streaming.classic.offset.migration.include.connector.name
仅在迁移模式为 strict 或 best_effort 时使用,取值必须与 v3 Connector 是否在通道名中包含 Connector 名称保持一致。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:高
- 必填:条件必填
behavior.on.null.values
有效值为 default、ignore。ignore 过滤 Kafka Tombstone;default 保留旧行为,将空 JSON 内容写入目标表。
- 类型:
STRING - 默认值:default
- 重要级别:低
- 必填:否
日志、指标与高级选项
jmx
是否启用 Connector 自定义 Snowflake 指标 MBean。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:高
- 必填:否
snowflake.streaming.client.provider.override.map
Snowpipe Streaming SDK 的高级覆盖项。仅应在 Snowflake Support 指导下使用。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
错误处理
errors.tolerance
有效值为 all、none。all 容忍 Connector 侧的摄取或校验记录错误,并可配合 DLQ 保留失败记录;none 使 Task 失败。
- 类型:
STRING - 默认值:none
- 重要级别:低
- 必填:否
errors.log.enable
记录已容忍的记录错误。启用详细框架消息日志前应评估记录中的敏感信息。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
errors.deadletterqueue.topic.name
在 errors.tolerance=all 时用于保存已容忍的失败记录;空值禁用 DLQ 输出。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
enable.mdc.logging
是否为 Connector 日志启用全局 MDC 上下文。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
enable.task.fail.on.authorization.errors
设为 true 时,已观察到的 Snowflake 授权错误会在 preCommit 阶段使 Task 失败。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
snowflake.compatibility.enable.autogenerated.table.name.sanitization
true 对自动生成的表名进行清理并转为大写,以兼容 v3;false 直接使用 Topic 名。特殊名称可使用显式 snowflake.topic2table.map。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:条件必填
snowflake.compatibility.enable.column.identifier.normalization
true 将列标识符规范化为大写,以兼容 v3。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:条件必填
snowflake.enable.schematization
true 将记录映射到独立列;false 写入兼容 v3 的 RECORD_CONTENT 和 RECORD_METADATA VARIANT 列。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:中
- 必填:条件必填
snowflake.autocreate.table.type
有效值为 snowflake、iceberg、none。前两者在表缺失时自动创建对应类型的表;none 在表缺失时失败。已有表始终按现有类型和结构使用。
- 类型:
STRING - 默认值:snowflake
- 重要级别:中
- 必填:否
snowflake.iceberg.create.table.options
追加到自动创建的托管 Iceberg 表 CREATE 语句中的 SQL 子句。不要包含 CATALOG、ENABLE_SCHEMA_EVOLUTION 或 ERROR_LOGGING;已有表会忽略该配置。
- 类型:
STRING - 默认值:空字符串
- 重要级别:低
- 必填:否
snowflake.cache.table.exists
是否缓存目标表存在性检查结果。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.cache.table.exists.expire.ms
表存在性缓存的过期时间,单位毫秒,最小值为 1。
- 类型:
LONG - 默认值:300000
- 重要级别:低
- 必填:否
snowflake.cache.pipe.exists
是否缓存 Pipe 存在性检查结果。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
snowflake.cache.pipe.exists.expire.ms
Pipe 存在性缓存的过期时间,单位毫秒,最小值为 1。
- 类型:
LONG - 默认值:300000
- 重要级别:低
- 必填:否
snowflake.topic2table.map.regex.replacement
true 允许在映射表名模板中使用 Java 正则捕获组替换;false 保留字面量 $ 和旧版不替换行为。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:低
- 必填:否
并发与转换
connector.class
使用 com.snowflake.kafka.connector.SnowflakeStreamingSinkConnector。这是 Kafka Connect 框架配置,不属于 Connector 自身的 ConfigDef。
- 类型:
STRING - 默认值:无固定默认值
- 重要级别:高
- 必填:是
tasks.max
Connector 请求的 Task 数量上限;实际有效并行度还受已分配 Topic 分区数限制。
- 类型:
INT - 默认值:1
- 重要级别:高
- 必填:否
tasks.max.enforce
设为 true 时,如果 Connector 返回的 Task 配置数超过 tasks.max,Kafka Connect 将其判为失败。
- 类型:
BOOLEAN - 默认值:true
- 重要级别:低
- 必填:否
- 已弃用:是
- 替代项:无
topics
以逗号分隔的 Topic 列表,与 topics.regex 互斥,两者必须且只能配置一个。
- 类型:
LIST - 默认值:空字符串
- 重要级别:高
- 必填:条件必填
topics.regex
用于订阅 Topic 的完整 Java 正则表达式,与 topics 互斥,两者必须且只能配置一个。
- 类型:
STRING - 默认值:空字符串
- 重要级别:高
- 必填:条件必填
key.converter
未设置时继承 Worker 的 Key Converter。Converter 专属子配置由所选插件决定,没有统一默认值。
- 类型:
CLASS - 默认值:
null(继承 Worker 配置) - 重要级别:低
- 必填:否
value.converter
未设置时继承 Worker 的 Value Converter。schemas.enable 等子配置由所选 Converter 决定,没有统一默认值。
- 类型:
CLASS - 默认值:
null(继承 Worker 配置) - 重要级别:低
- 必填:否
errors.deadletterqueue.topic.replication.factor
仅在 Kafka Connect 自动创建缺失的 DLQ Topic 时使用,取值必须适合目标 Kafka 集群的 broker 数量。
- 类型:
SHORT - 默认值:3
- 重要级别:中
- 必填:否
errors.deadletterqueue.context.headers.enable
为框架写入的 DLQ 记录添加 __connect.errors.* 上下文 Header。
- 类型:
BOOLEAN - 默认值:false
- 重要级别:中
- 必填:否
最佳实践
将多个 Topic 路由到明确的目标表
适用业务场景:一个 Connector 需要把不同 Topic 写入不同 Snowflake 表,或 Topic 名称不能直接作为目标表名。 配置示例:在基础配置中加入显式映射。容忍记录错误并保留失败记录
适用业务场景:个别记录可能在 Converter 或 Connector 侧处理失败,但不能停止持续写入任务,同时需要保存失败记录以便排查。 配置示例:启用 Connect 容错和 DLQ,且 DLQ Topic 不得被当前订阅匹配。3,broker 数不足时应预创建 Topic 或显式调整 errors.deadletterqueue.topic.replication.factor。
监控
监控内容
监控 Kafka Connect Worker 健康状态、Connector 和 Task 状态、吞吐、延迟、offset 提交、错误与重试,以及 Worker JVM 的堆内存、GC、线程和 CPU;启用错误容忍和 DLQ 后再关注 DLQ 写入量、失败记录增长和投递错误。导入 Grafana 大盘
下载 AutoMQ Connect Cluster Dashboard,在 Grafana 中选择已采集 Kafka Connect 指标的数据源,确认指标标签与大盘变量匹配后导入 JSON 文件。限制条件
- 不提供跨 Task、Topic 或 Kafka 分区的全局记录顺序。
- Snowflake 通道已提交 offset 加一是单个 Topic 分区的恢复和 Kafka 提交边界;无法读取通道已提交状态时,该分区不会被报告为可安全提交。
- 已知已提交 offset 的重放记录会被跳过,但该机制不构成跨 Task、分区或表的原子 exactly-once 保证。
snowflake.autocreate.table.type=none时缺失目标表会导致任务初始化失败。- 托管 Iceberg 表不兼容
snowflake.validation=client_side。 - offset 提交和恢复以 Topic 分区为边界,不提供跨多个分区或表的原子检查点。
- 后端限流、通道恢复失败和不可恢复摄取错误可能导致任务失败。
常见问题
任务启动时提示必填配置缺失怎么办?
检查 Snowflake URL、用户、数据库、Schema 和角色。默认认证方式是snowflake_jwt,还需提供 snowflake.private.key;使用 OAuth 时改为检查 OAuth 客户端 ID 和密钥。
为什么配置了 Topic 后仍然没有写入目标表?
确认topics 和 topics.regex 只配置一个,并检查订阅是否匹配实际 Topic。再检查数据库、Schema、角色权限以及显式映射和自动建表设置。
为什么启动时出现 Classic 兼容性配置错误?
默认兼容性校验会要求snowflake.validation=client_side、两个 v3 规范化开关为 true,并显式设置 schematization 和 Classic offset migration。新建 4.1.0 Connector 且不需要 v3 兼容时,可将 snowflake.streaming.validate.compatibility.with.classic 设为 false。从 v3 迁移时,选择 strict 或 best_effort,并根据旧 Connector 是否启用了 snowflake.streaming.channel.name.include.connector.name 设置 snowflake.streaming.classic.offset.migration.include.connector.name。
为什么结构化数据无法转换?
启用 schematization 时,Value 必须能转换为 Map 或 Struct,且不能使用 StringConverter 或 ByteArrayConverter。检查 Converter、消息结构和目标表权限。为什么错误记录没有出现在 DLQ?
确认errors.tolerance=all、DLQ Topic 已设置且未被当前订阅匹配,并检查 Worker 是否提供 ErrantRecordReporter 以及 Kafka 是否允许写入该 Topic。