如果你需要任何帮助,请在 repository 中提交 issue,或在 ClickHouse public Slack 中提出问题。
许可证
环境要求
版本兼容性矩阵
主要特性
- 开箱即用,具备精确一次语义。其底层依托 ClickHouse 一项名为 KeeperMap 的全新核心特性 (由连接器用作状态存储) ,从而实现极简架构。
- 支持 3rd-party 状态存储:当前默认使用内存,也可使用 KeeperMap (即将支持 Redis) 。
- 核心集成:由 ClickHouse 构建、维护并提供支持。
- 持续针对 ClickHouse Cloud 进行测试。
- 支持按声明的 schema 和无 schema 方式插入数据。
- 支持 ClickHouse 的所有数据类型。
安装说明
准备连接信息
你的 ClickHouse Cloud 服务的连接信息可在 ClickHouse Cloud 控制台中查看。
选择一个服务,然后点击 Connect:
选择 HTTPS。连接信息会显示在示例
curl 命令中。
如果你使用的是自管理 ClickHouse,则连接信息由你的 ClickHouse 管理员配置。
常规安装说明
- 从 ClickHouse Kafka Connect Sink repository 的 Releases 页面下载包含连接器 JAR 文件的 ZIP 压缩包。
- 解压 ZIP 文件,并将其内容复制到所需位置。
- 在 Connect 属性文件的 plugin.path 配置中添加插件目录路径,以便 Confluent Platform 能够找到该插件。
- 在配置中提供 topic 名称、ClickHouse 实例的 hostname 和密码。
- 重启 Confluent Platform。
- 如果您使用 Confluent Platform,请登录 Confluent Control Center UI,确认 ClickHouse Sink 是否出现在可用连接器列表中。
配置选项
- 连接信息:hostname (必填) 和端口 (可选)
- 用户凭据:密码 (必填) 和用户名 (可选)
- 连接器类:
com.clickhouse.kafka.connect.ClickHouseSinkConnector(必填) - topics 或 topics.regex:要轮询的 Kafka topic——topic 名称必须与表名一致 (必填)
- key 和 value 转换器:根据 topic 中的数据类型进行设置。如果尚未在工作线程配置中定义,则为必填。
目标表
预处理
支持的数据类型
-
(1) - 仅当 ClickHouse 设置中启用了
input_format_binary_read_json_as_string=1时,才支持 JSON。这仅适用于 RowBinary 格式家族,并且该设置会影响 insert 请求中的所有列,因此这些列都应为字符串。在这种情况下,连接器 会将 STRUCT 转换为 JSON 字符串。 -
(2) - 当 struct 包含
oneof这类 union 时,应将 converter 配置为不要为 field 名称添加 prefix/suffix。ProtobufConverter提供了generate.index.for.unions=false设置。
配置范例
基础版配置
localhost:8443 上运行了启用 SSL 的 ClickHouse server,数据采用无 schema 的 JSON 格式。
上述连接器配置要求你在工作线程配置中通过
connector.client.config.override.policy=All 启用客户端覆盖。更多信息请参阅 Kafka Connect 文档。多个 topic 的基本配置
含 DLQ 的基本配置
配合不同数据格式使用
支持 Avro schema
Avro 类型映射
io.confluent.connect.avro.AvroConverter 定义,它是 Kafka Connect 官方的 Avro 序列化器/反序列化器实现。有关转换逻辑的更多高级信息,请参阅 Kafka Connect 文档。
✅:支持
❌:不支持
️⚠️:部分支持
有关 Kafka Connect 类型与 ClickHouse 类型之间的映射,请参阅支持的数据类型。
不支持的 Avro schema
- 固定长度的
decimal逻辑类型
- Nullable 联合类型
- 记录中的联合类型
支持 Protobuf schema
Protobuf 类型映射
io.confluent.connect.protobuf.ProtobufConverter 定义,它是 Kafka Connect 官方的 Protobuf 序列化/反序列化实现。有关转换逻辑的更多高级信息,请参阅 Kafka Connect 文档。
✅:支持
❌:不支持
️⚠️:部分支持
有关 Kafka Connect 类型与 ClickHouse 类型之间的映射,请参阅支持的数据类型。
关于将 oneof 字段映射为 ClickHouse 列的说明
oneof) 映射为 ClickHouse 的 Variant 类型。请改为在 ClickHouse 表 schema 中,将 oneof 字段列为单独的可空字段。
例如:
不支持的 Protobuf schema
- 多消息联合类型 (在 CH 版本 26.1 之前)
allow_experimental_nullable_tuple_type=1 时支持此 schema (参见此文档页面) 。
支持 JSON schema
支持 String
内部缓冲
poll() 调用返回的记录,并将其作为更大的批次刷新到 ClickHouse。在每次轮询都会为各个分区产生大量小批次的工作负载中,这可以提升吞吐量。
关键行为:
bufferCount控制刷新前缓冲的记录数。bufferFlushTime设置刷新缓冲记录前的最长等待时间 (以毫秒为单位) 。bufferFlushTime仅在bufferCount > 0时生效。bufferCount=0和bufferFlushTime=0会使缓冲保持禁用状态 (默认行为) 。- 当
exactlyOnce=true时,不支持缓冲。
exactlyOnce=false 禁用 exactly-once 模式,或者使用 bufferCount=0 禁用缓冲。
示例:
日志
监控
ClickHouse 特有指标
Kafka 生产者/消费者指标
records-sent-total:发送到 topic 的记录总数bytes-sent-total:发送到 topic 的总字节数record-send-rate:每秒发送记录的平均速率byte-rate:每秒发送字节的平均速率compression-rate:实现的压缩率
records-sent-total:发送到分区的记录总数bytes-sent-total:发送到分区的总字节数records-lag:分区当前的滞后量records-lead:分区当前的超前量replica-fetch-lag:副本的滞后信息
connection-creation-total:与 Kafka 节点建立的连接总数connection-close-total:关闭的连接总数request-total:发送到节点的请求总数response-total:从节点接收的响应总数request-rate:每秒请求的平均速率response-rate:每秒响应的平均速率
- 吞吐量:跟踪数据摄取速率
- 滞后:识别瓶颈和处理延迟
- 压缩:衡量数据压缩效率
- 连接健康状况:监控网络连通性和稳定性
Kafka Connect Framework 指标
task-count:连接器中的任务总数running-task-count:当前正在运行的任务数量paused-task-count:当前已暂停的任务数量failed-task-count:已失败的任务数量destroyed-task-count:已销毁的任务数量unassigned-task-count:未分配的任务数量
running、paused、failed、destroyed、unassigned
错误指标:
deadletterqueue-produce-failures:DLQ 写入失败次数deadletterqueue-produce-requests:DLQ 写入尝试总次数last-error-timestamp:最近一次错误的时间戳records-skip-total:因错误而跳过的记录总数records-retry-total:重试的记录总数errors-total:发生的错误总数
offset-commit-failures:偏移量 提交失败次数offset-commit-avg-time-ms:偏移量 提交的平均耗时offset-commit-max-time-ms:偏移量 提交的最长耗时put-batch-avg-time-ms:处理一个批次的平均耗时put-batch-max-time-ms:处理一个批次的最长耗时source-record-poll-total:轮询到的记录总数
监控最佳实践
- 监控消费者滞后:按分区跟踪
records-lag,以识别处理瓶颈 - 跟踪错误率:关注
errors-total和records-skip-total,以发现数据质量问题 - 关注任务健康状态:监控任务状态指标,确保任务正常运行
- 衡量吞吐量:使用
records-send-rate和byte-rate跟踪摄取性能 - 监控连接状态:检查节点级连接指标,以发现网络问题
- 跟踪压缩效率:使用
compression-rate优化数据传输
限制
- 不支持删除。
- 批次大小继承自 Kafka 消费者属性。
- 使用 KeeperMap 实现 exactly-once 时,如果偏移量发生变化或被回退,需要删除 KeeperMap 中该特定 topic 的内容。 (更多详情请参阅下方的故障排查指南)
性能调优与吞吐量优化
何时需要进行性能调优?
- 高吞吐量工作负载:当需要从 Kafka topic 中以每秒数百万条的速度处理事件时
- 滞后:当连接器无法跟上数据生成速度,导致滞后持续增加时
- 资源受限:当你需要优化 CPU、内存或网络资源的使用时
- 多个 topic:当需要同时消费多个高流量 topic 时
- 消息较小:当需要处理大量小消息,而这些消息适合通过服务端批处理来提升效率时
- 处理的数据量较低或中等 (< 10,000 条消息/秒)
- 对你的使用场景来说,消费者滞后稳定且可接受
- 默认连接器设置已经满足你的吞吐量要求
- 你的 ClickHouse 集群可以轻松处理传入负载
理解数据流
- Kafka Connect Framework 在后台从 Kafka topic 拉取消息
- 连接器轮询 框架内部 buffer 中的消息
- 连接器分批 根据轮询大小将消息组成批次
- ClickHouse 接收 通过 HTTP/S 发送的批次 insert
- ClickHouse 处理 该 insert (同步或异步)
Kafka Connect 批次大小调优
拉取设置
fetch.min.bytes:框架将值传给 连接器 之前所需的最小数据量 (默认值:1 字节)fetch.max.bytes:单个请求可拉取的最大数据量 (默认值:52428800 / 50 MB)fetch.max.wait.ms:如果未达到fetch.min.bytes,返回数据前最长等待的时间 (默认值:500 毫秒)
在 Confluent Cloud 上,如需调整这些设置,需要通过 Confluent Cloud 提交支持工单。
轮询设置
max.poll.records:单次轮询返回的最大记录数 (默认值:500)max.partition.fetch.bytes:每个分区可拉取的最大数据量 (默认值:1048576 / 1 MB)
在 Confluent Cloud 上,调整这些设置需要通过 Confluent Cloud 提交支持工单。
高吞吐量场景下的推荐设置
上述属性要求你在 worker 配置中通过
connector.client.config.override.policy=All 启用客户端配置覆盖。更多信息请参阅 Kafka Connect 文档。- 更大的批次 = 更好的 ClickHouse 摄取性能、更少的 parts、更低的开销
- 更大的批次 = 更高的内存使用量,以及端到端延迟可能增加
- 批次过大 = 可能会导致 timeout、OutOfMemory 错误,或超出
max.poll.interval.ms
异步插入
何时使用异步插入
- 大量小批次:你的连接器会频繁发送小批次 (每批少于 1000 行)
- 高并发:多个连接器任务同时向同一张表写入数据
- 分布式部署:在不同主机上运行多个连接器实例
- parts 创建开销:你遇到了“parts 过多”错误
- 混合工作负载:同时处理实时摄取和查询工作负载
- 你已经以可控频率发送大批次 (每批超过 10,000 行)
- 你需要数据立即可见 (查询必须立刻看到数据)
- 使用
wait_for_async_insert=0时,精确一次语义与你的需求冲突 - 你的使用场景更适合通过客户端侧批处理优化来改进
异步插入的工作原理
- 从连接器接收插入查询
- 将数据写入内存缓冲区 (而不是立即写入磁盘)
- 向连接器返回成功 (如果
wait_for_async_insert=0) - 在满足以下任一条件时,将缓冲区刷写到磁盘:
- 缓冲区达到
async_insert_max_data_size(默认值:100 MB) - 自首次插入起已过去
async_insert_busy_timeout_ms毫秒 (默认值:1000 ms) - 累积查询数达到上限 (
async_insert_max_query_number,默认值:100)
- 缓冲区达到
启用异步插入
clickhouseSettings 配置参数中:
async_insert=1:启用异步插入wait_for_async_insert=1(推荐) :连接器会等数据刷写到 ClickHouse 存储后再确认。可提供交付保障。wait_for_async_insert=0:连接器在数据缓冲后立即确认。性能更好,但如果服务器在刷写前崩溃,数据可能会丢失。
调优 异步插入 行为
async_insert_max_data_size(默认值:104857600 / 100 MB) :触发 flush 前的最大缓冲区大小async_insert_busy_timeout_ms(默认值:1000) :触发 flush 前的最长等待时间 (毫秒)async_insert_stale_timeout_ms(默认值:0) :距离上次 insert 后触发 flush 的时间 (毫秒)async_insert_max_query_number(默认值:100) :触发 flush 前的最大查询数
- 优点:parts 更少、merge 性能更好、CPU 开销更低,并且在高并发下吞吐量更高
- 注意事项:数据无法立即被查询到,端到端延迟会略有增加
- 风险:如果
wait_for_async_insert=0,server 崩溃时可能导致数据丢失;缓冲区过大时可能带来内存压力
具有精确一次语义的异步插入
exactlyOnce=true 时:
wait_for_async_insert=1 与 exactly-once 配合使用,以确保只有在数据持久化后才会提交偏移量。
有关异步插入的更多信息,请参阅 ClickHouse 异步插入文档。
连接器并行度
每个连接器的任务数
- 有效任务数的上限 = topic 分区数
- 每个任务都会维护自己与 ClickHouse 的连接
- 任务越多,开销越高,还可能导致资源争用
tasks.max 设为与 topic 分区数相同,然后再根据 CPU 和吞吐量指标进行调整。
批处理时忽略分区
exactlyOnce=false 时使用。此设置可通过创建更大的批次来提高吞吐量,但会失去分区内的顺序保证。
多个高吞吐量 topic
topic2TableMap 将 topic 映射到表,并且在插入时遇到瓶颈,导致消费滞后,则建议改为每个 topic 单独创建一个 连接器。
出现这种情况的主要原因是,目前批次会串行插入到各个表中。
建议:对于多个高流量 topic,建议为每个 topic 部署一个独立的 连接器 实例,以最大化并行插入吞吐量。
ClickHouse 表引擎注意事项
MergeTree:最适合大多数场景,兼顾查询和 insert 性能ReplicatedMergeTree:高可用性所必需,但会增加复制开销*MergeTree配合适当的ORDER BY:针对你的查询模式进行优化
连接池与超时
socket_timeout(默认值:30000 毫秒) :读取操作的最长等待时间connection_timeout(默认值:10000 毫秒) :建立连接的最长等待时间
监控与性能问题排查
- 滞后:使用 Kafka 监控工具跟踪各分区的滞后
- 连接器指标:通过 JMX 监控
receivedRecords、recordProcessingTime、taskProcessingTime(参见 Monitoring) - ClickHouse 指标:
system.asynchronous_inserts:监控异步插入缓冲区的使用情况system.parts:监控 parts 数量,以发现合并问题system.merges:监控正在进行的合并system.events:跟踪InsertedRows、InsertedBytes、FailedInsertQuery
最佳实践总结
- 先使用默认配置,然后根据实际性能进行测量和调优
- 尽量使用更大的批次:如有可能,目标是每次插入 10,000-100,000 行
- 在发送大量小批次或高并发场景下,使用异步插入
- 使用精确一次语义时,始终启用
wait_for_async_insert=1 - 水平扩展:将
tasks.max提高到最多与分区数相同 - 每个高流量 topic 使用一个连接器,以获得最大吞吐量
- 持续监控:跟踪滞后、part 数量和合并活动
- 充分测试:在生产部署前,务必在接近真实的负载下测试配置更改
示例:高吞吐量配置
上述连接器配置要求你在 worker 配置中通过
connector.client.config.override.policy=All 启用客户端覆盖。更多信息请参阅 Kafka Connect 文档。- 每次轮询最多处理 10,000 条记录
- 跨分区合并批次,以实现更大规模的插入
- 使用 16 MB 缓冲区的异步插入
- 运行 8 个并行任务 (与分区数保持一致)
- 针对吞吐量进行了优化,而非严格顺序
故障排查
”topic [someTopic] 的分区 [0] 状态不一致”
此调整可能会对 exactly-once 语义产生影响。
“连接器会重试哪些错误?”
ClickHouseException- 这是 ClickHouse 可能抛出的通用异常。 它通常会在 server 过载时抛出,以下错误代码尤其被视为暂时性错误:- 3 - UNEXPECTED_END_OF_FILE
- 107 - FILE_DOESNT_EXIST
- 159 - TIMEOUT_EXCEEDED
- 164 - READONLY
- 202 - TOO_MANY_SIMULTANEOUS_QUERIES
- 203 - NO_FREE_CONNECTION
- 209 - SOCKET_TIMEOUT
- 210 - NETWORK_ERROR
- 241 - MEMORY_LIMIT_EXCEEDED
- 242 - TABLE_IS_READ_ONLY
- 252 - TOO_MANY_PARTS
- 285 - TOO_FEW_LIVE_REPLICAS
- 319 - UNKNOWN_STATUS_OF_INSERT
- 425 - SYSTEM_ERROR
- 999 - KEEPER_EXCEPTION
SocketTimeoutException- 当套接字超时时会抛出此异常。UnknownHostException- 当主机无法解析时会抛出此异常。IOException- 当网络出现问题时会抛出此异常。
“我的所有数据都是空的/全是 0”
_ 作为分隔符) 。这样一来,表中的字段就会采用 “field1_field2_field3” 这种格式 (即 “before_id”、“after_id” 等) 。
“我想在 ClickHouse 中使用 Kafka 键”
KeyToValue 转换将键移到 value 字段中 (作为新的 _key 字段) :