Skip to main content

概述

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 同名的表。
将占位符替换为实际资源,端点使用 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
  • 重要级别:低
  • 有效值 / 注意事项:索引类型不区分大小写,可使用 PRIMARYCOLUMNSTOREUNIQUESHARDKEY,可追加键名;值为逗号分隔列名。只在创建缺失表时生效,不修改已有表,并对所有自动创建的表采用同一组键配置。

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
  • 重要级别:低
  • 有效值 / 注意事项truefalse。只有目标表存在 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
  • 重要级别:低
  • 有效值 / 注意事项:不区分大小写,可使用 GZipLZ4SkipSkip 表示不压缩;驱动和服务端必须支持对应的文件扩展与传输方式。

重试与指标标签

max.retries

设置发生 SQL 异常后的 Connector 级最大重试次数。
  • 类型int
  • 默认值10
  • 重要级别:中
  • 有效值 / 注意事项:必须大于等于 00 表示首次 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。 配置示例
替换占位符,并确认输入 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,在 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.retriesretry.backoff.ms;路由、Schema 访问、序列化、指标注册和本地文件流等 Task 内部错误不会使用该重试预算,也不会由 Connector 写入 Kafka Connect DLQ。
  • Connector 不为单次写入提供独立的记录数或字节数批次配置;有效批次由 Kafka Connect Consumer Poll 和路由结果决定,多张目标表在一个 Task 内串行写入。

常见问题

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

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

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

自动建表要求记录值带 Connect Schema,并要求数据库账号具备表探测和建表权限;无 Schema Map 或其他无 Schema 值必须使用预创建表。对于已有表,Connector 不执行 ALTER TABLE,因此应检查 Value Converter 产生的字段、fields.whitelistfields.blacklistsinglestore.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,但这些机制都有适用边界。应使用业务唯一标识审计重复数据,并把链路按可能重放进行设计。