概述
ClickHouse Sink Connector 从 Kafka Topic 消费记录,将记录值写入 ClickHouse 数据表,适用于行为事件、应用日志和业务流水的实时分析链路。通常一个 Topic 对应一张同名表,也可以显式映射到已有业务表。结构化记录和 JSON 对象的字段按名称与目标列关联;CSV/TSV 字符串则按目标表输入列顺序写入。 Connector 根据值 Converter 生成的 Connect 数据模型处理带 Schema 的结构化记录、无 Schema 的 JSON 对象或字符串记录。写入以追加行为为基础,不会把普通 Kafka 消息自动转换为按键更新或删除。可选的 Debezium CDC 处理将变更事件转换为版本化行,配合目标表引擎表达当前状态。前置条件
- ClickHouse 服务端需为 23.3 或更高版本。
- 预先创建目标数据库和数据表,列名、类型及缺失字段的默认值规则应与输入记录匹配。
- ClickHouse 账号需能读取目标表元数据并执行 INSERT。
- 启用自动加列时,账号还需具备目标表 ALTER 权限。
- 启用
exactlyOnce时,需配置可用的 KeeperMap 状态存储及其建表、读写权限,并确认目标表引擎与去重窗口满足重放去重要求。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka Topic 和 ClickHouse 数据库及目标表,确认网络连通与访问权限;资源准备和 Connector 管理操作参见管理 Connector。下面采用无 Schema 的 JSON 对象输入,通过 HTTPS 写入与 Topic 同名的表。hostname 不含协议或端口,HTTPS 端口默认是 8443;若服务端使用其他端口,增加 port。HTTP 部署需同时设置 ssl=false 和对应 HTTP 端口,例如 8123,协议切换不会自动修改端口。通过安全的凭据管理方式提供密码,不将实际凭据提交到配置仓库。
例如,输入 JSON 对象包含整数 event_id 和字符串 event_type 时,在上述 ClickHouse 数据库中执行以下建表语句,并将 <events-table> 替换为与 <events-topic> 完全相同的名称。
{"event_id":1,"event_type":"page_view"},不添加 Connect 的 schema / payload 包装。此配置未启用持久状态去重;故障恢复或重放可能产生重复行,普通 MergeTree 不会根据 ORDER BY 自动去重。
配置
以下默认值为配置定义的声明值;个别参数省略时的实际行为在注意事项中单独说明。列表类型通常使用逗号分隔,dateTimeFormats 的多字段格式使用分号。
Connector 身份与输入订阅
connector.class
指定 Connector 实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:使用
com.clickhouse.kafka.connect.ClickHouseSinkConnector。
tasks.max
允许创建的最大 Task 数。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:大于等于
1;实际有效并行度受输入分区数量和分配限制,增加 Task 不保证线性提升吞吐。
topics
指定消费的 Topic 列表。
- 类型:
list - 默认值:空列表
- 重要级别:高
- 必填:与
topics.regex二选一 - 有效值 / 注意事项:逗号分隔;与
topics.regex必须且只能设置一个非空值,不可包含自身 DLQ Topic。
topics.regex
通过正则表达式订阅 Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:高
- 必填:与
topics二选一 - 有效值 / 注意事项:使用 Java 正则表达式,不可匹配自身 DLQ Topic;新匹配 Topic 也需准备目标表。
ClickHouse 连接与认证
hostname
ClickHouse 主机名。
- 类型:
string - 默认值:无
- 重要级别:高
- 必填:是
- 有效值 / 注意事项:非空,不包含协议或端口。
port
ClickHouse HTTP 或 HTTPS 接口端口。
- 类型:
int - 默认值:
8443 - 重要级别:高
- 有效值 / 注意事项:
1至65535;不是原生 TCP 接口端口,也不会随ssl自动切换。
username
ClickHouse 用户名。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:省略参数时连接实现回退为
default;显式填写账号,避免依赖省略与空值的不同处理。
password
ClickHouse 账号密码。
- 类型:
password - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:读取时去除首尾空格;不要使用依赖首尾空格的密码,不在日志或公开示例中填写实际凭据。
ssl
选择是否通过 HTTPS 连接。
- 类型:
boolean - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false;省略时连接实现可能回退为 HTTP,建议显式设置并与port匹配。
client_version
选择 ClickHouse 客户端实现。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:建议明确填写
V1或V2。只有精确的V1选择旧客户端,其他值包括显式空字符串都选择 V2;省略时回退为 V1。
jdbcConnectionProperties
为连接 URL 添加查询属性。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:使用
key=value&key=value格式,可省略前导?。连接 URL 可能被记录到日志,不通过此项传递密码或令牌。
clickhouseSettings
为写入请求指定 ClickHouse 服务端设置。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:逗号分隔的
key=value,每项只包含一个=。连接实现逐项补入用户未指定的input_format_skip_unknown_fields=1、wait_end_of_query=1、async_insert=0、send_progress_in_http_headers=1,即使本项已包含其他设置也会补入,不覆盖用户值;这些不属于本项声明默认值。不要覆盖 Connector 管理的insert_deduplication_token。异步插入需等待实际入库,wait_for_async_insert=0不提供持久化确认。
ssl_socket_sni
覆盖 TLS 握手中的 SNI 主机名。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:仅在特殊 TLS 路由需求下使用;V1 客户端设置非空值时还会关闭证书及主机名验证,不能将其作为无安全影响的 SNI 调整。
代理连接
proxyType
指定 ClickHouse 连接代理类型。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:推荐值为
IGNORE、DIRECT、HTTP、SOCKS;空字符串或无法识别的值忽略代理。V2 对非IGNORE类型统一构建 HTTP 代理,不保留 SOCKS 或 DIRECT 的原有语义。
proxyHost
代理服务器主机名。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:仅在非
IGNORE代理类型下使用,启用代理时填写实际主机名。
proxyPort
代理服务器端口。
- 类型:
int - 默认值:
-1 - 重要级别:低
- 有效值 / 注意事项:启用代理时显式填写有效网络端口,默认值不表示可用端口。
数据库与表路由
database
默认目标数据库。
- 类型:
string - 默认值:
default - 重要级别:低
- 有效值 / 注意事项:数据库需预先存在;启用 Topic 拆分时可由 Topic 名中的数据库部分覆盖。
topic2TableMap
将 Topic 映射到目标表。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:逗号分隔的
topic=table,键值去除首尾空格,重复键以后项为准。未映射 Topic 使用同名表;映射值是表名,不是跨库database.table路由语法,也不会创建表。
enableDbTopicSplit
从 Topic 名拆分数据库与表路由名称。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:启用时显式设置
dbTopicSplitChar;仅恰好拆成两段时覆盖数据库,随后按拆出的 Topic 部分查找topic2TableMap,其他情况保留原数据库与 Topic 名。
dbTopicSplitChar
数据库与 Topic 部分的分隔符。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:仅在
enableDbTopicSplit=true时使用;推荐_、-或.,按字面分隔符处理而非正则表达式。
suppressTableExistenceException
抑制目标表不存在时的异常。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:启用后缺表记录可能被忽略且 Offset 继续推进;不创建表,也不是等待目标表出现的机制。要求完整写入的链路应保持关闭。
值格式与字段处理
value.converter
将 Kafka 消息值反序列化为 Connect 数据。
- 类型:
class - 默认值:
null - 重要级别:低
- 有效值 / 注意事项:
null表示使用 Worker 的 Converter。类需实现 Converter;带 Schema 的 Struct、无 Schema 的 Map 和 String 使用不同写入路径,Avro 或 Protobuf 需另行提供相应 Converter 及依赖。
customInsertFormat
用于字符串输入的自定义格式选项。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:此项影响配置界面的格式推荐值,不单独控制实际写入路径;字符串的实际格式由
insertFormat决定。
insertFormat
指定 String 值的插入格式。
- 类型:
string - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:字符串输入应明确选择
CSV、TSV或JSON,其中 JSON 使用 JSONEachRow。NONE无法写入字符串;不要依赖未知值回退到 JSON 的行为。CSV/TSV 字段顺序需符合目标表输入顺序。
bypassRowBinary
将带 Schema 记录从二进制写入切换为 JSONEachRow。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:不改变 Converter,也不会自动补齐字段;默认带 Schema 路径使用 RowBinary 或 RowBinaryWithDefaults。
dateTimeFormats
为日期时间字段指定解析模式。
- 类型:
list - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:实际输入使用分号分隔多个
field=pattern,例如created_at=yyyy-MM-dd HH:mm:ss;updated_at=yyyy-MM-dd HH:mm:ss。模式遵循 Java DateTimeFormatter 语法,非法模式会导致配置读取失败。
bypassFieldCleanup
保留 JSON 写入路径中不属于目标表的顶层字段。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:关闭时按目标列清理额外顶层字段;开启后仍受服务端
input_format_skip_unknown_fields和类型检查约束。
debeziumCDCEnabled
将 Debezium Envelope 转为版本化目标行。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:要求值为 Struct 且 Schema 名以
.Envelope结尾;目标表使用ReplacingMergeTree(_version, is_deleted),并按业务主键设计排序键。删除事件(op=d)使用before并写入is_deleted=1;创建、快照读取和更新事件(op=c/r/u)使用after并写入is_deleted=0,同时添加_version。不支持通过 truncate 事件(op=t)清空目标表;删除标记也不是执行 SQL DELETE,不保证立即物理删除。
表结构与自动加列
tableRefreshInterval
周期刷新目标表元数据。
- 类型:
long - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:
0至600,单位秒;0不启用周期刷新。已有表主要在列数增加时重新读取,不保证跟踪列类型修改、删列或重命名。
bypassSchemaValidation
跳过记录与目标表的 Schema 对齐检查。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:仅跳过 Connector 的检查,不绕过序列化约束和 ClickHouse 类型校验,不作为不兼容数据的修复方式。
auto.evolve
按输入 Schema 为目标表添加缺失列。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:面向带真实 Connect Schema 的结构化记录使用,不依赖无 Schema JSON 或 String 自动推导业务类型。需要 ALTER 权限;不创建业务表,不修改已有列类型或删除列。按写入批次最后一条记录检查字段,应保持批次 Schema 演进顺序一致。
auto.evolve.ddl.refresh.retries
自动加列后等待新列元数据可见的重试次数。
- 类型:
int - 默认值:
3 - 重要级别:低
- 有效值 / 注意事项:大于等于
0;只控制 DDL 后的元数据刷新,不是写入请求或框架错误的统一重试次数。
auto.evolve.struct.to.json
自动加列时将 Connect STRUCT 推导为 ClickHouse JSON 列。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅作用于自动演进类型映射,服务端需支持 JSON 类型及相关设置。使用二进制写入 JSON 列时,还需在
clickhouseSettings设置input_format_binary_read_json_as_string=1或true。
clusterName
为自动演进 DDL 指定 ClickHouse 集群。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:非空时添加
ON CLUSTER,空字符串仅执行本地 DDL;需有效集群及 DDL 权限。仅填写受信任集群名,不填写 SQL 片段;与keeperOnCluster独立。
批量与内部缓冲
ignorePartitionsWhenBatching
将同 Topic 不同分区的记录合并写入。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:用于
exactlyOnce=false;exactlyOnce=true且无缓冲时不采用此选项,启用缓冲时设置为true会导致 Task 启动失败。不提供跨分区全局顺序。
bufferCount
跨多次消费累积记录的数量阈值。
- 类型:
int - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:大于等于
0;0不启用内部缓冲。非 exactly-once 模式达到总量阈值后刷新,阈值不是最大批次上限。exactly-once 模式按 Topic/分区发送固定大小的完整批次,要求bufferFlushTime=0和ignorePartitionsWhenBatching=false,不足数量的尾部继续等待。
bufferFlushTime
内部缓冲的时间触发阈值。
- 类型:
long - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:大于等于
0,单位毫秒;0禁用时间触发,仅bufferCount>0时有效。时间条件在 Task 处理调用时检查,不是独立定时刷新承诺;exactly-once 缓冲模式必须保持0。
状态、去重与 Offset
exactlyOnce
启用基于 KeeperMap 的持久状态与批次去重路径。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:开启后按 Topic/分区记录处理前后状态。重放去重依赖可用的 KeeperMap、目标引擎去重支持及有效窗口、稳定数据与批次边界、可靠的入库确认;单独设置开关不构成无条件 exactly-once 保证。数据写入与状态写入不是同一事务,不提供跨表或跨分区原子性。
zkPath
KeeperMap 状态存储路径。
- 类型:
string - 默认值:
/kafka-connect - 重要级别:低
- 有效值 / 注意事项:非空白且以
/开头;用于exactlyOnce=true的状态存储,应规划状态隔离。
zkDatabase
KeeperMap 状态表名称。
- 类型:
string - 默认值:
connect_state - 重要级别:低
- 有效值 / 注意事项:虽名为 Database,实际表示状态表标识符,不是单独的数据库创建配置;使用有效且受信任的表名。
keeperOnCluster
指定 Keeper 状态表所用的 ClickHouse 集群。
- 类型:
string - 默认值:空字符串
- 重要级别:低
- 有效值 / 注意事项:用于自托管集群的持久状态场景,不等同于业务表 DDL 的
clusterName。
tolerateStateMismatch
容忍部分批次范围与已有状态不匹配的情况。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:某些落后于已有状态的记录可能被跳过;不是 Offset 回退后的重新入库开关,也不保证修复所有状态错误。开启前需评估数据完整性。
reportInsertedOffsets
让直接写入策略在提交前报告已处理写入范围的 Offset。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:内部缓冲策略自行维护已刷新范围,不由本项控制;非 exactly-once 跨分区合批仍使用 Worker 当前 Offset。错误容忍及缺表抑制可使“处理完成”不同于“实际入库”,此项不是逐行写入成功证明。
超时与错误处理
timeoutSeconds
ClickHouse 驱动操作超时。
- 类型:
int - 默认值:
30 - 重要级别:低
- 有效值 / 注意事项:
0至600,单位秒;与插入响应等待上限clickhouseClientInsertTimeoutMs分开配置。
retryCount
ClickHouse 客户端辅助操作重试次数。
- 类型:
int - 默认值:
3 - 重要级别:低
- 有效值 / 注意事项:
3至10;影响连接探测及查询等辅助操作,不是所有 INSERT 的统一最大重试次数。
clickhouseClientInsertTimeoutMs
V2 客户端单次插入响应等待上限。
- 类型:
long - 默认值:
240000 - 重要级别:低
- 有效值 / 注意事项:大于等于
100,单位毫秒;V1 插入不使用此独立等待上限。应小于消费者有效max.poll.interval.ms并留出处理和重试余量;超时或取消等待不能证明服务端未接受数据。
consumer.override.max.poll.interval.ms
覆盖 Sink 消费者相邻 poll 的最大间隔。
- 类型:
int - 默认值:无 Connector 级独立默认值;底层客户端默认
300000 - 重要级别:中
- 有效值 / 注意事项:大于等于
1,单位毫秒;Worker 消费者配置可改变有效值,是否允许覆盖取决于 Worker 的客户端覆盖策略。应大于插入等待时间并留余量,避免慢写入期间发生再均衡。
errors.retry.timeout
Kafka Connect 框架可重试错误的总重试时长。
- 类型:
long - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:单位毫秒,
0不重试,-1无限重试;只作用于相应框架可重试处理路径,不是所有 ClickHouse 写入错误的自动重试保证,不替代retryCount。
errors.tolerance
控制错误容忍策略。
- 类型:
string - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:
none或all;同名框架项重要级别为中。all可跳过被容忍的错误而继续推进,不是重试开关;插件写入失败可能按整个分组报告,不能视为逐行精确隔离坏记录。没有可用 DLQ 时,被跳过的数据不等于已持久保存。
errors.deadletterqueue.topic.name
为错误报告指定死信队列 Topic。
- 类型:
string - 默认值:空字符串
- 重要级别:中
- 有效值 / 注意事项:空字符串不记录 DLQ;需配合错误容忍策略及 Worker 错误报告机制,且不可被本 Connector 订阅。配置 Topic 名不等于报告已持久成功,应检查实际 DLQ 写入及其错误。
最佳实践
多条业务事件流写入已有分析表
适用业务场景:订单事件和支付事件使用不同 Kafka Topic,而 ClickHouse 已有稳定命名的分析表,希望保留现有表名而不要求 Topic 与表同名。 配置示例:在快速开始配置中覆盖topics,增加以下映射。两个表都位于快速开始指定的数据库中,需预先创建并分别匹配各自事件结构。
汇集零散事件以减少小批次写入
适用业务场景:行为或应用事件持续到达,但每次消费只有少量记录,希望累积后写入 ClickHouse,减少频繁小批次插入,同时让低流量记录在后续处理调用中也有机会刷新。 配置示例:在快速开始配置中添加以下设置。该方案保持非 exactly-once 模式,不与持久状态去重缓冲方案混用。1000 毫秒不是严格入库延迟上界,也没有独立后台计时器保证准点入库。单次消费的记录数可能超过数量阈值,最终写入仍按 Topic/分区分组。根据事件大小、Worker 内存和可接受延迟调整数量与时间,并检查实际可见延迟;此方案不消除重放重复,不能据此承诺故障恢复 exactly-once。
监控
监控内容
关注 Kafka Connect 集群健康、Connector 与 Task 的运行和失败状态、消费及写入吞吐、消费积压与端到端延迟、Offset 提交进度和失败、错误与重试变化,以及 Worker JVM 的堆内存、GC 和线程信号;启用错误容忍与 DLQ 报告后,还需关注 DLQ 活动和报告失败,避免把 Task 持续运行误认为全部数据已入库。导入 Grafana 大盘
下载共享的 Kafka Connect Grafana Dashboard,确保 Kafka Connect 指标已采集到 Prometheus 兼容数据源,且集群、Worker、Connector 和 Task 等标签与大盘查询一致;在 Grafana 中导入 JSON 文件并选择对应数据源。限制条件
- Connector 不自动创建业务数据库或数据表,自动演进仅添加缺失列。
- Kafka tombstone 不触发目标表删除,常规写入不执行按键 SQL UPDATE 或 DELETE。
- 输入值须转换为受支持的 Struct、Map、String 或 null,不支持任意 Connect 原始标量直接写入。
- 不提供跨分区或跨 Task 的全局顺序,也不保证目标查询结果按写入顺序排列。
- 数据与 Keeper 状态写入不是同一事务,不提供跨表、跨分区原子事务。
- exactly-once 内部缓冲不支持时间触发刷新和跨分区合批,不足数量的尾部不会在提交、关闭或停止时强制写出。
- 自动加列不修改已有列类型或删除列,周期元数据刷新也不是任意表结构变更的自动适配机制。
常见问题
Task 无法连接 ClickHouse 或启动时找不到表
检查hostname 是否误含协议、port 是否为 HTTP/HTTPS 端口,以及 ssl 是否与服务端匹配。显式填写 ssl 和 client_version,避免省略参数的回退行为影响连接。确认服务端满足版本要求、账号能读取表元数据,且 database 下已有目标表;使用映射时核对 Topic 和表名。不要通过启用 suppressTableExistenceException 掩盖缺表问题,应先创建或修正目标表。
输入 JSON 后出现格式或字段类型错误
确认消息是普通 JSON 对象还是带schema / payload 的 Connect JSON:普通对象应使用快速开始的 JsonConverter 并关闭 Schema 包装解析。核对目标列名、类型以及缺失字段的 DEFAULT 或 Nullable 定义;StringConverter 输入还需明确设置 insertFormat。不要只打开 bypassSchemaValidation,它不会让 ClickHouse 接受不兼容类型。
开启缓冲后少量数据迟迟不可见
检查记录数量是否达到bufferCount。非 exactly-once 模式可以设置非零 bufferFlushTime,但时间条件需等到 Task 的后续处理调用才检查;若业务要求尽快可见,可关闭内部缓冲。exactly-once 缓冲仅写完整固定大小批次,不支持时间刷新,应减小数量阈值或关闭缓冲,不要添加冲突的时间触发设置。
重启后出现重复行
写入成功与 Kafka Offset 提交之间存在故障窗口,重放可能再次插入;普通 MergeTree 的排序键不提供唯一约束。先检查提交失败、再均衡和超时,确认服务端是否已接受请求。需要持久状态去重时,评估 KeeperMap、目标引擎去重条件和窗口、稳定批次及入库确认,不能只设置exactlyOnce=true 或自定义去重 token 就认定不会重复。
回退 Offset 后出现状态不匹配
持久状态仍保留已处理范围,外部 Offset 回退不会自动协调 KeeperMap。先停止写入,核对所需重放范围与对应 Topic/分区状态,再制定状态清理或独立重放方案并评估重复风险。tolerateStateMismatch=true 可能跳过旧范围,不是重新导入历史数据的修复开关。
Task 仍在运行,但数据缺失或 DLQ 没有记录
检查是否启用errors.tolerance=all 或缺表异常抑制,两者都可能让未入库记录按已处理推进。确认 DLQ Topic 不在订阅范围内,Worker 错误报告机制可用,并检查报告权限、发送错误和实际消息;配置了 Topic 名也不代表报告必定持久成功。要求完整写入时保持失败可见,修复根因后按已确认的 Offset 范围安排重放。