> ## 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.

# SingleStore Sink Connector

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

## 概述

SingleStore Sink Connector 消费 Kafka Topic 中的记录，并通过 `LOAD DATA LOCAL INFILE` 将记录值批量写入 SingleStore 数据表。它位于 Kafka 事件流与用于实时分析或事务处理的 SingleStore 数据库之间，适合持续汇入业务事件、应用日志和状态更新数据。

对于单 Topic 输入，Connector 默认使用 Topic 名作为目标表名，也可以把该 Topic 映射到固定表；启用记录字段路由后，还可以根据字段值把同一批记录分组写入多张表。带 Connect Schema 的 Struct 字段通常映射为同名列，带 Connect Schema 的单值记录写入 `data` 列；无 Schema 的 Map 也可写入已存在且结构兼容的表。Connector 可根据带 Connect Schema 的记录创建缺失表，但不会为已有表执行 Schema 演进。

## 前置条件

* 目标数据库必须已经存在；数据库账号需能查询表元数据并执行 `LOAD DATA LOCAL INFILE`，使用自动建表时还需具备建表权限，保留默认元数据功能时还需能创建和读写元数据表。
* SingleStore 服务端和 JDBC 连接必须允许 `LOAD DATA LOCAL INFILE`，并支持所选的 GZip、LZ4 或不压缩传输方式。
* 无 Schema 记录需要预先创建目标表；已有表的列名、类型、可空性及主键或唯一键应能接收 Connector 生成的行，后续 Schema 变化需在数据库侧维护。

## 授权许可

使用 Apache License 2.0。

## 快速开始

提前准备 Connect Cluster、Kafka Topic、SingleStore 数据库和具有相应权限的数据库账号，并确认网络连通和访问权限。具体准备和管理操作请参阅 [管理 Connector](../manage-connectors)。下面的最小实用配置通过 SingleStore Helios 云工作区端点读取字符串记录，并写入与 Topic 同名的表。

```properties theme={null}
connector.class=com.singlestore.kafka.SingleStoreSinkConnector
topics=<input-topic>
value.converter=org.apache.kafka.connect.storage.StringConverter
connection.clientEndpoint=<singlestore-endpoint>
connection.database=<database-name>
connection.user=<database-user>
connection.password=<database-password>
```

将占位符替换为实际资源，端点使用 JDBC 驱动接受的 `host:port` 形式。预先创建与 Topic 同名、包含可接收字符串的 `data` 列的目标表；通过安全的凭据管理方式提供密码，不将真实凭据提交到版本控制系统。

## 配置

### Connector 身份与输入订阅

#### `connector.class`

选择 SingleStore Sink Connector 实现类。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **必填**：是
* **有效值 / 注意事项**：使用 `com.singlestore.kafka.SingleStoreSinkConnector`。

#### `tasks.max`

设置 Kafka Connect 最多可创建的 Sink Task 数量。

* **类型**：`int`
* **默认值**：`1`
* **重要级别**：高
* **有效值 / 注意事项**：必须至少为 `1`。实际有效并行度受输入 Topic 分区数和分区分配限制；多个 Task 不提供跨 Task 的写入顺序或同键串行化。

#### `topics`

指定 Connector 消费的 Kafka Topic 列表。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：高
* **必填**：与 `topics.regex` 二选一
* **有效值 / 注意事项**：使用逗号分隔 Topic 名称；必须与非空的 `topics.regex` 二选一，两者同时设置或同时为空都会校验失败。

#### `topics.regex`

通过 Java 正则表达式选择 Connector 消费的 Kafka Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **必填**：与 `topics` 二选一
* **有效值 / 注意事项**：必须是合法且非空的 Java 正则表达式，并与 `topics` 互斥。新增匹配 Topic 时同时确认其目标表路由和表结构。

### 记录转换

#### `key.converter`

覆盖该 Connector 从 Worker 继承的 Kafka Key Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：`null` 表示继承 Worker 配置；显式值必须是具有公共无参构造方法的 Converter 类。Connector 不使用 Kafka Key 生成目标列或行标识。

#### `value.converter`

覆盖该 Connector 从 Worker 继承的 Kafka Value Converter。

* **类型**：`class`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：`null` 表示继承 Worker 配置；显式值必须是具有公共无参构造方法的 Converter 类。Converter 产生的 Schema 和值决定自动建表、字段映射和行序列化行为。

### SingleStore 连接与认证

#### `connection.ddlEndpoint`

设置自管理部署中用于表探测、建表和查询的 DDL 端点。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **必填**：与 `connection.clientEndpoint` 二选一
* **有效值 / 注意事项**：与 `connection.clientEndpoint` 互斥。未设置非空 `connection.dmlEndpoints` 时，该端点也用于 DML 写入；地址格式和连通性由 JDBC 驱动校验。

#### `connection.clientEndpoint`

设置 SingleStore Helios 云工作区用于 DDL 和 DML 的单一端点。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：高
* **必填**：与 `connection.ddlEndpoint` 二选一
* **有效值 / 注意事项**：与 `connection.ddlEndpoint` 及非空的 `connection.dmlEndpoints` 互斥；选用后所有数据库操作使用该端点。

#### `connection.database`

设置写入数据和元数据的 SingleStore 数据库。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **必填**：是
* **有效值 / 注意事项**：数据库名称会加入 JDBC URL。配置定义不拒绝空字符串，但数据库必须实际存在且账号可访问，否则 Task 无法启动。

#### `connection.user`

设置 JDBC 连接使用的 SingleStore 用户。

* **类型**：`string`
* **默认值**：`root`
* **重要级别**：高
* **有效值 / 注意事项**：使用满足表探测、建表和写入需求的最小权限账号；凭据会在 Task 启动建立 JDBC 连接时验证。

#### `connection.password`

设置 SingleStore 用户的密码。

* **类型**：`password`
* **默认值**：`null`
* **重要级别**：高
* **有效值 / 注意事项**：使用安全的配置注入方式提供真实密码，不在日志、文档或版本控制系统中暴露凭据。

#### `connection.dmlEndpoints`

设置用于 DML 写入的 SingleStore Aggregator 端点列表。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：使用逗号分隔端点，仅能与 `connection.ddlEndpoint` 配合；与 `connection.clientEndpoint` 冲突。省略或设置为空列表时，DML 使用已选择的 DDL 或 Client 端点。

#### `params.<value>`

向 SingleStore JDBC 驱动传递动态连接参数。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：将实际参数写成 `params.<driver-property>=<value>`，Connector 会移除 `params.` 前缀后传给驱动。参数名和值遵循随 Connector 打包的 JDBC 驱动规则；部分参数可能包含密钥或证书密码。不要把 `params.allowLocalInfile` 设为会禁止本地文件流的值。

### 表结构与字段映射

#### `tableKey.<index_type>[.<name>]`

为 Connector 自动创建的表添加键定义。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：索引类型不区分大小写，可使用 `PRIMARY`、`COLUMNSTORE`、`UNIQUE`、`SHARD` 或 `KEY`，可追加键名；值为逗号分隔列名。只在创建缺失表时生效，不修改已有表，并对所有自动创建的表采用同一组键配置。

#### `fields.whitelist`

只保留列出的顶层记录字段参与后续路由和写入。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：使用逗号分隔且区分大小写的字段名。适用于 Struct 和无 Schema Map；`null` 或空列表表示不限制。与黑名单同时使用时先应用白名单。

#### `fields.blacklist`

从记录中排除列出的顶层字段。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：使用逗号分隔且区分大小写的字段名。适用于 Struct 和无 Schema Map；`null` 或空列表表示不排除。与白名单同时包含同一字段时，黑名单最终移除该字段。

#### `singlestore.columnToField.<tableName>.<columnName>`

把目标表列映射到 Kafka 记录中的字段路径。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：属性名中的表名和列名必须各占一个点分段，不能包含额外的点；属性值可使用点分隔的嵌套字段路径。映射只对解析出的目标表名完全匹配时生效。写入已有表时，缺失路径会产生 SQL `NULL`；如果 Connector 需要按 Schema 自动建表但无法解析映射字段的 Schema，建表会失败。

### 目标表路由

#### `singlestore.tableName.<topicName>`

把一个 Kafka Topic 映射到固定的 SingleStore 表。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：使用实际 Topic 名替换 `<topicName>`，值填写目标表名；没有匹配映射时使用 Topic 名作为表名。不能与 `singlestore.recordToTable.mappingField` 同时使用。

#### `singlestore.recordToTable.mappingField`

指定用于逐条记录选择目标表的字段路径。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：可使用点分隔的 Struct 字段或 Map Key 路径，并需配置对应的 `singlestore.recordToTable.mapping.<value>`。不能与任何 Topic 到表映射同时使用；字段路径不存在、字段值为 `null` 或没有匹配映射的记录会被跳过。

#### `singlestore.recordToTable.mapping.<value>`

把一个路由字段值映射到 SingleStore 表。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：使用实际路由值替换 `<value>`，值填写目标表名。至少一个具体映射要求同时设置 `singlestore.recordToTable.mappingField`；只有命中映射的记录会写入。

### 写入、去重与压缩

#### `singlestore.filter`

为生成的 `LOAD DATA` 语句添加 `WHERE` 过滤表达式。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：表达式会直接拼接到 SQL，不经过 Connector 解析或参数化。只使用受信任且适用于所有目标表的表达式，避免把外部输入直接写入该配置。

#### `singlestore.upsert`

为 `LOAD DATA` 启用 `REPLACE` 行为。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：`true` 或 `false`。只有目标表存在 PRIMARY 或 UNIQUE 键冲突时才会替换已有行；这是整行替换，不是部分字段合并，也不构成 Kafka Offset 与数据库之间的事务保证。

#### `singlestore.metadata.allow`

启用元数据表和基于批次首条记录标识的重复批次抑制。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：启用时，Connector 创建或使用元数据表，在写入前检查批次标识，并将元数据与该批次的表写入放在同一数据库事务中。该机制只覆盖相同首条 Kafka 坐标的批次重放，不能视为端到端 exactly-once。

#### `singlestore.metadata.table`

设置 Connector 使用的元数据表名称。

* **类型**：`string`
* **默认值**：`kafka_connect_transaction_metadata`
* **重要级别**：低
* **有效值 / 注意事项**：在 `singlestore.metadata.allow=true` 时生效。使用当前数据库中可创建、查询和写入的表名；配置不校验标识符是否合法。

#### `singlestore.loadDataCompression`

选择 `LOAD DATA LOCAL INFILE` 数据流的压缩方式。

* **类型**：`string`
* **默认值**：`GZip`
* **重要级别**：低
* **有效值 / 注意事项**：不区分大小写，可使用 `GZip`、`LZ4` 或 `Skip`。`Skip` 表示不压缩；驱动和服务端必须支持对应的文件扩展与传输方式。

### 重试与指标标签

#### `max.retries`

设置发生 SQL 异常后的 Connector 级最大重试次数。

* **类型**：`int`
* **默认值**：`10`
* **重要级别**：中
* **有效值 / 注意事项**：必须大于等于 `0`；`0` 表示首次 SQL 异常即失败。重试预算在一次成功写入后重置，且与 Kafka Connect 通用错误处理配置相互独立。

#### `retry.backoff.ms`

设置 SQL 写入重试前请求的等待时间。

* **类型**：`int`
* **默认值**：`3000`
* **重要级别**：中
* **有效值 / 注意事项**：必须大于等于 `0`，单位毫秒；只在 `max.retries` 仍有剩余次数时使用，没有指数退避或抖动。

#### `custom.metric.tags`

为每个 Task 的 JMX ObjectName 添加自定义标签。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：使用逗号分隔的 `key=value`；每项必须恰好包含一个 `=`，重复 Key 以后项为准。值会进行 JMX 清理，但 Key 按原样加入 ObjectName，应避免非法 JMX 字符。

## 最佳实践

### 将业务 Topic 映射到稳定的目标表名

**适用业务场景**：首次接入时，Kafka Topic 名包含环境、版本或组织前缀，但数据库需要使用稳定、简洁的业务表名，并希望后续调整 Topic 命名时保持目标表不变。

**配置示例**：

```properties theme={null}
connector.class=com.singlestore.kafka.SingleStoreSinkConnector
topics=orders.events.v1
value.converter=org.apache.kafka.connect.storage.StringConverter
connection.clientEndpoint=<singlestore-endpoint>
connection.database=<database-name>
connection.user=<database-user>
connection.password=<database-password>
singlestore.tableName.orders.events.v1=orders_events
```

将连接占位符替换为实际资源，并预先创建 `orders_events`，或授予账号根据输入 Schema 自动创建该表的权限。

**关键说明**：显式映射把 Kafka 资源命名与数据库表命名解耦，便于维护权限、Schema 和下游查询。Topic 映射与记录字段路由互斥；需要按记录内容分表时应改用 `singlestore.recordToTable.*`，并为所有有效路由值配置目标表。

### 用目标表业务键替换重复状态行

**适用业务场景**：目标表用于保存客户、订单或设备的当前状态，同一业务实体的更新已通过 Kafka 分区保持所需顺序，但重放或重复事件可能再次到达，希望依据数据库 PRIMARY 或 UNIQUE 键替换已有整行。

**配置示例**：

```properties theme={null}
connector.class=com.singlestore.kafka.SingleStoreSinkConnector
topics=customer-updates
value.converter=org.apache.kafka.connect.json.JsonConverter
connection.clientEndpoint=<singlestore-endpoint>
connection.database=<database-name>
connection.user=<database-user>
connection.password=<database-password>
singlestore.tableName.customer-updates=customers
singlestore.upsert=true
```

将连接占位符替换为实际资源，并预先创建 `customers` 表及能唯一标识业务实体的 PRIMARY 或 UNIQUE 键。输入记录使用 `JsonConverter` 接受的带 Schema JSON 表示，字段需与目标列兼容，并包含数据库键对应的字段。

**关键说明**：`singlestore.upsert=true` 使用 SingleStore `LOAD DATA REPLACE`，冲突时替换整行而非只更新变化字段。替换顺序取决于记录的处理顺序，不会按事件时间自动选择最新版本；同一业务键跨分区写入时也没有跨 Task 协调。该设置可以降低重复写入的影响，但 Kafka Offset 提交与数据库事务不原子，元数据表也只按批次首条记录抑制部分重放，因此不能据此宣称 exactly-once。

### 按 Kafka 分区逐步增加写入并行度

**适用业务场景**：Connector 已稳定运行，输入 Topic 有多个分区且消费 Lag 持续增长，希望增加并行 Task，让多个分区同时写入 SingleStore。

**配置示例**：

```properties theme={null}
connector.class=com.singlestore.kafka.SingleStoreSinkConnector
topics=<input-topic>
tasks.max=3
value.converter=org.apache.kafka.connect.storage.StringConverter
connection.clientEndpoint=<singlestore-endpoint>
connection.database=<database-name>
connection.user=<database-user>
connection.password=<database-password>
```

替换占位符，并确认输入 Topic 至少有多个可分配分区。先小幅提高 `tasks.max`，同时观察消费 Lag、Task 写入延迟、SQL 重试、数据库负载和 Worker JVM 资源。

**关键说明**：`tasks.max` 是上限，增加到高于可分配分区数不会产生更多有效工作。不同 Task 之间没有全局顺序或同键协调；同一业务键可能由多个分区并发写入时，应先明确分区键、数据库键和 `REPLACE` 语义。

## 监控

### 监控内容

关注 Kafka Connect Worker 健康状态、Connector 和 Task 状态、输入吞吐、消费 Lag、处理延迟、Offset 提交、错误、SQL 重试以及 Worker JVM 的 CPU、内存和垃圾回收信号；同时核对 SingleStore 实际写入行数和数据库错误。仅在部署启用了相应 Kafka Connect 错误处理时关注 DLQ 活动，Connector 内部路由、序列化和写入错误不一定进入 DLQ。

### 导入 Grafana 大盘

确认 Kafka Connect 指标已接入 Grafana 数据源，且采集标签满足大盘筛选条件；下载 [Kafka Connect Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，在 Grafana 中导入 JSON 并选择对应数据源。

## 限制条件

* SingleStore 数据写入与 Kafka Offset 提交不属于同一事务；数据库提交后、Offset 提交前发生故障可能重放记录，默认元数据表只按批次首条记录标识抑制部分重复，不能提供端到端 exactly-once。
* 无 Schema 的记录不能用于自动创建缺失表；已有表不会执行自动 Schema 演进，新增、删除、重命名或更改类型的字段需先在数据库侧处理。
* Kafka Record Key 和 Header 不会写入目标列，也不会自动成为表主键或更新标识；Kafka Tombstone 不会删除目标行。
* `singlestore.upsert=true` 依赖目标表 PRIMARY 或 UNIQUE 键，并执行整行 `REPLACE`，不是部分字段合并。
* 记录字段路由只写入命中 `singlestore.recordToTable.mapping.<value>` 的记录；字段路径不存在、字段值为 `null` 或未配置映射的记录会被跳过，且没有专门的跳过记录指标或 DLQ 报告。
* Connector 只对 SQL 异常应用 `max.retries` 和 `retry.backoff.ms`；路由、Schema 访问、序列化、指标注册和本地文件流等 Task 内部错误不会使用该重试预算，也不会由 Connector 写入 Kafka Connect DLQ。
* Connector 不为单次写入提供独立的记录数或字节数批次配置；有效批次由 Kafka Connect Consumer Poll 和路由结果决定，多张目标表在一个 Task 内串行写入。

## 常见问题

### Task 启动时提示端点配置冲突或无法连接 SingleStore，如何处理？

确认 `connection.ddlEndpoint` 和 `connection.clientEndpoint` 只设置一个；使用 `connection.clientEndpoint` 时移除 `connection.dmlEndpoints`。随后检查 `connection.database` 是否存在、账号密码是否正确，以及账号能否从所有运行 Task 的 Worker 访问并查询目标数据库。自管理部署需要多个写入端点时，保留 `connection.ddlEndpoint`，再配置逗号分隔的 `connection.dmlEndpoints`。

### 为什么目标表没有自动创建，或者写入提示列不匹配？

自动建表要求记录值带 Connect Schema，并要求数据库账号具备表探测和建表权限；无 Schema Map 或其他无 Schema 值必须使用预创建表。对于已有表，Connector 不执行 `ALTER TABLE`，因此应检查 Value Converter 产生的字段、`fields.whitelist`、`fields.blacklist` 和 `singlestore.columnToField.*` 是否与目标列名、类型和可空性一致，再先完成数据库 Schema 变更。

### 为什么启用记录字段路由后部分消息没有写入？

检查 `singlestore.recordToTable.mappingField` 的点分字段路径是否能在 Struct 或 Map 中解析，并确认每个实际字段值都有对应的 `singlestore.recordToTable.mapping.<value>`。路径不存在、字段值为 `null` 或值未映射时记录会被跳过；还要确认没有同时配置冲突的 `singlestore.tableName.<topicName>`。

### SQL 错误为什么会重复出现，最终 Task 仍然失败？

Connector 会对所有 SQL 异常重复整个写入批次，最多重试 `max.retries` 次，每次重试前等待 `retry.backoff.ms`；这些重试用尽后，如果同一批次再次发生 SQL 异常，Task 会失败。检查 Worker 日志和 SingleStore 错误，区分连接中断等临时问题与权限、SQL 表达式、表结构或重复键等永久问题；修正根因后再恢复 Task，不要仅通过扩大重试次数掩盖不可恢复错误。

### 为什么故障恢复后出现重复行？

数据库写入成功后，Kafka Offset 仍需由 Kafka Connect 单独提交；两者之间发生故障会使同一批记录再次投递。保持 `singlestore.metadata.allow=true` 可抑制首条 Kafka 坐标相同的批次重放，保存最新状态的表还可使用稳定 PRIMARY 或 UNIQUE 键配合 `singlestore.upsert=true`，但这些机制都有适用边界。应使用业务唯一标识审计重复数据，并把链路按可能重放进行设计。
