> ## Documentation Index
> Fetch the complete documentation index at: https://docs.automq.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Snowflake Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Snowflake Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

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](../manage-connectors)。

```properties theme={null}
connector.class=com.snowflake.kafka.connector.SnowflakeStreamingSinkConnector
tasks.max=1
topics=orders
snowflake.url.name=https://<account-url>
snowflake.user.name=<snowflake-user>
snowflake.private.key=<rsa-private-key>
snowflake.database.name=<database-name>
snowflake.schema.name=<schema-name>
snowflake.role.name=<role-name>
snowflake.authenticator=snowflake_jwt
snowflake.streaming.validate.compatibility.with.classic=false
snowflake.enable.schematization=false
snowflake.autocreate.table.type=snowflake
value.converter=org.apache.kafka.connect.json.JsonConverter
```

将尖括号中的值替换为实际值。示例订阅 `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 名称不能直接作为目标表名。

配置示例：在基础配置中加入显式映射。

```properties theme={null}
connector.class=com.snowflake.kafka.connector.SnowflakeStreamingSinkConnector
tasks.max=2
topics=orders,customers
snowflake.url.name=https://<account-url>
snowflake.user.name=<snowflake-user>
snowflake.private.key=<rsa-private-key>
snowflake.database.name=<database-name>
snowflake.schema.name=<schema-name>
snowflake.role.name=<role-name>
snowflake.authenticator=snowflake_jwt
snowflake.streaming.validate.compatibility.with.classic=false
snowflake.enable.schematization=false
snowflake.topic2table.map=orders:ORDERS,customers:CUSTOMERS
value.converter=org.apache.kafka.connect.json.JsonConverter
```

关键说明：精确 Topic 匹配优先，重复或重叠映射会被拒绝；未加引号的表名转为大写。

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

适用业务场景：个别记录可能在 Converter 或 Connector 侧处理失败，但不能停止持续写入任务，同时需要保存失败记录以便排查。

配置示例：启用 Connect 容错和 DLQ，且 DLQ Topic 不得被当前订阅匹配。

```properties theme={null}
connector.class=com.snowflake.kafka.connector.SnowflakeStreamingSinkConnector
tasks.max=1
topics=orders
snowflake.url.name=https://<account-url>
snowflake.user.name=<snowflake-user>
snowflake.private.key=<rsa-private-key>
snowflake.database.name=<database-name>
snowflake.schema.name=<schema-name>
snowflake.role.name=<role-name>
snowflake.authenticator=snowflake_jwt
snowflake.streaming.validate.compatibility.with.classic=false
snowflake.enable.schematization=false
errors.tolerance=all
errors.deadletterqueue.topic.name=orders-dlq
errors.deadletterqueue.context.headers.enable=true
value.converter=org.apache.kafka.connect.json.JsonConverter
```

关键说明：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](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 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。
