概述
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。下面的最小实用配置通过 SingleStore Helios 云工作区端点读取字符串记录,并写入与 Topic 同名的表。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 命名时保持目标表不变。 配置示例:orders_events,或授予账号根据输入 Schema 自动创建该表的权限。
关键说明:显式映射把 Kafka 资源命名与数据库表命名解耦,便于维护权限、Schema 和下游查询。Topic 映射与记录字段路由互斥;需要按记录内容分表时应改用 singlestore.recordToTable.*,并为所有有效路由值配置目标表。
用目标表业务键替换重复状态行
适用业务场景:目标表用于保存客户、订单或设备的当前状态,同一业务实体的更新已通过 Kafka 分区保持所需顺序,但重放或重复事件可能再次到达,希望依据数据库 PRIMARY 或 UNIQUE 键替换已有整行。 配置示例:customers 表及能唯一标识业务实体的 PRIMARY 或 UNIQUE 键。输入记录使用 JsonConverter 接受的带 Schema JSON 表示,字段需与目标列兼容,并包含数据库键对应的字段。
关键说明:singlestore.upsert=true 使用 SingleStore LOAD DATA REPLACE,冲突时替换整行而非只更新变化字段。替换顺序取决于记录的处理顺序,不会按事件时间自动选择最新版本;同一业务键跨分区写入时也没有跨 Task 协调。该设置可以降低重复写入的影响,但 Kafka Offset 提交与数据库事务不原子,元数据表也只按批次首条记录抑制部分重放,因此不能据此宣称 exactly-once。
按 Kafka 分区逐步增加写入并行度
适用业务场景:Connector 已稳定运行,输入 Topic 有多个分区且消费 Lag 持续增长,希望增加并行 Task,让多个分区同时写入 SingleStore。 配置示例: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,在 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,但这些机制都有适用边界。应使用业务唯一标识审计重复数据,并把链路按可能重放进行设计。