Skip to main content

概述

Snowflake Sink Connector 将 Kafka Topic 中的记录持续写入 Snowflake 表,位于 Kafka 与 Snowflake 之间的数据写入链路中。它使用 Snowpipe Streaming 将每个 Topic 分区的数据写入对应的 Snowflake 摄取通道,可根据 Topic 名称自动确定目标表,也可通过映射将 Topic 路由到指定表。 Connector 可以把结构化记录展开为 Snowflake 列,也可以将记录写入 RECORD_CONTENTRECORD_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_CONTENTRECORD_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_jwtoauth
  • 类型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.schematizationsnowflake.streaming.classic.offset.migration;迁移模式为 strictbest_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_sideclient_sideserver_side 要求目标表启用错误日志;托管 Iceberg 表不支持 client_side
  • 类型STRING
  • 默认值:server_side
  • 重要级别:高
  • 必填:否

snowflake.streaming.classic.offset.migration

有效值为 skipbest_effortstrict。仅当 SSv2 通道没有已提交 offset 时才尝试读取 v3 Classic 通道:strict 在旧通道不存在时失败,best_effort 回退到 Kafka consumer group offset,skip 不读取旧通道。
  • 类型STRING
  • 默认值:skip
  • 重要级别:高
  • 必填:条件必填

snowflake.streaming.classic.offset.migration.include.connector.name

仅在迁移模式为 strictbest_effort 时使用,取值必须与 v3 Connector 是否在通道名中包含 Connector 名称保持一致。
  • 类型BOOLEAN
  • 默认值:false
  • 重要级别:高
  • 必填:条件必填

behavior.on.null.values

有效值为 defaultignoreignore 过滤 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

有效值为 allnoneall 容忍 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_CONTENTRECORD_METADATA VARIANT 列。
  • 类型BOOLEAN
  • 默认值:true
  • 重要级别:中
  • 必填:条件必填

snowflake.autocreate.table.type

有效值为 snowflakeicebergnone。前两者在表缺失时自动创建对应类型的表;none 在表缺失时失败。已有表始终按现有类型和结构使用。
  • 类型STRING
  • 默认值:snowflake
  • 重要级别:中
  • 必填:否

snowflake.iceberg.create.table.options

追加到自动创建的托管 Iceberg 表 CREATE 语句中的 SQL 子句。不要包含 CATALOGENABLE_SCHEMA_EVOLUTIONERROR_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 名称不能直接作为目标表名。 配置示例:在基础配置中加入显式映射。
关键说明:精确 Topic 匹配优先,重复或重叠映射会被拒绝;未加引号的表名转为大写。

容忍记录错误并保留失败记录

适用业务场景:个别记录可能在 Converter 或 Connector 侧处理失败,但不能停止持续写入任务,同时需要保存失败记录以便排查。 配置示例:启用 Connect 容错和 DLQ,且 DLQ Topic 不得被当前订阅匹配。
关键说明:DLQ 处理 Kafka Connect 和 Connector 侧可容忍的记录错误,不接收 Snowflake 服务端校验失败的记录;后端限流、通道恢复失败和不可恢复摄取错误仍可能导致任务失败。Kafka Connect 自动创建 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 后仍然没有写入目标表?

确认 topicstopics.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 迁移时,选择 strictbest_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。