尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

ClickHouse Kafka Connect Sink 连接器实战:将 Kafka 数据实时写入 ClickHouse

ClickHouse Kafka Connect Sink 连接器实战:将 Kafka 数据实时写入 ClickHouse ClickHouse Kafka Connect Sink 连接器实战将 Kafka 数据实时写入 ClickHouse【免费下载链接】lagoOpen Source Metering and Usage Based Billing API ⭐️ Consumption tracking, Subscription management, Pricing iterations, Payment orchestration Revenue analytics项目地址: https://gitcode.com/GitHub_Trending/la/lago导读本文围绕 Lago 开源计费平台仓库中随附的官方ClickHouse Kafka Connect Sink 连接器v1.3.4展开系统讲解如何借助 Kafka Connect 将 Kafka本项目使用 RedpandaTopic 中的事件数据持续投递到 ClickHouse 表并重点剖析官方提供的KeyToValueSingle Message TransformationSMT——它能把 Kafka 消息 Key 原样落为 ClickHouse 表中的一列。读完本文你将掌握该连接器的功能边界、交付保证、部署方式含本仓库 docker-compose 集成示例、建表与连接器配置的全套实战方法以及官方性能测试基准的使用思路。一、连接器是什么Kafka Topic 到 ClickHouse 表的管道根据连接器官方 READMEextra/kafka-connect/clickhouse-kafka-connect-v1.3.4/doc/README.md的定位clickhouse-kafka-connect 是 ClickHouse 官方出品的 Kafka Connect Sink 连接器职责非常单一而明确——把 Kafka Topic 中的数据实时投递到 ClickHouse 表中。在典型的流式数据链路中它的位置如下Kafka / Redpanda Topic │ ▼ Kafka ConnectSink 连接器本主题 │ ▼ ClickHouse 表MergeTree 家族引擎这一点与 Lago 仓库的数据架构高度契合仓库的 docker-compose.dev.yml 中事件数据先写入 KafkaRedpanda的events_raw、events_enriched、events_enriched_expanded等主题见 scripts/create-topics.sh 相关主题创建逻辑随后由 events-processor 消费处理而extra/kafka-connect/目录正是仓库为 Kafka Connect 插件准备的挂载目录ClickHouse 则作为分析侧存储clickhouse服务使用clickhouse/clickhouse-server:26.2-alpine镜像。也就是说本仓库已经为Kafka → ClickHouse的数据管道预留了完整的运行环境本文讲解的连接器正是这条管道中的核心 Sink 组件。连接器版本与基本属性仓库中随附的插件清单文件 manifest.json 给出了连接器的权威元信息汇总如下属性值名称clickhouse-kafka-connect版本v1.3.4组件类型sink交付保证exactly_once精确一次支持的编码any任意格式配合 converter 使用SMT 支持single_message_transforms: trueKafka Connect API支持Confluent Control Center 集成支持运行要求ClickHouse v22.5 及以上许可证Apache License 2.0连接器的可执行产物是 lib/clickhouse-kafka-connect-v1.3.4-confluent.jarConfluent 发行版 jar连同manifest.json一起被 docker-compose.dev.yml 挂载进 Kafka Connect 容器。二、设计要点Exactly-Once 交付语义官方 README 明确提示该连接器的完整设计与 Exactly-Once 交付语义的实现原理由单独的设计文档design document专门阐述。关于 Exactly-Once 需要理解两点这是连接器宣称的交付保证manifest.json 中delivery_guarantee: [exactly_once]与之呼应。在 Sink 场景下Exactly-Once 通常依赖 Kafka 的事务机制与下游写入的幂等性配合实现涉及 offset 管理与目标表写入的事务性协调。代价与取舍精确一次语义通常会带来吞吐与延迟上的开销生产环境应根据业务对不重不漏的要求在at-least-once与exactly-once之间权衡。ClickHouse 本身是分析型数据库通常配合ReplacingMergeTree等去重引擎在读取侧兜底。说明设计文档位于连接器官方仓库中本仓库Lago随附的只有 README 与 jar 产物未包含设计文档正文如需深入了解内部实现请查阅官方项目仓库中的设计文档。三、在本仓库中部署 Kafka Connect 与 ClickHouse 连接器本仓库的 docker-compose.dev.yml 已经内置了完整的 Kafka Connect 运行环境可以直接观察它是如何接住这个连接器的。3.1 Kafka Connect 容器redpanda-kafka-connectredpanda-kafka-connect: image: docker.redpanda.com/redpandadata/connectors:latest hostname: redpanda-kafka-connect volumes: - ./extra/kafka-connect:/opt/kafka/connect-plugins ports: - 8083:8083 environment: CONNECT_CONFIGURATION: | key.converterorg.apache.kafka.connect.converters.ByteArrayConverter value.converterorg.apache.kafka.connect.converters.ByteArrayConverter group.idredpanda-kafka-connect-group offset.storage.topic_connectors_offsets config.storage.topic_connectors_configs status.storage.topic_connectors_status config.storage.replication.factor-1 offset.storage.replication.factor-1 status.storage.replication.factor-1 CONNECT_BOOTSTRAP_SERVERS: redpanda:9092 CONNECT_GC_LOG_ENABLED: true CONNECT_HEAP_OPTS: -Xms1G -Xmx1G CONNECT_PLUGIN_PATH: /opt/kafka/connect-plugins这段配置揭示了几个关键部署事实插件发现机制./extra/kafka-connect目录即clickhouse-kafka-connect-v1.3.4/与debezium-connector-postgres/所在目录被挂载到容器的/opt/kafka/connect-plugins并通过CONNECT_PLUGIN_PATH声明为插件搜索路径。Kafka Connect 会扫描该目录下的 jar 包与manifest.json自动识别出 ClickHouse Sink 连接器与 Debezium CDC 连接器。REST API 管理面容器暴露8083:8083这是 Kafka Connect 的标准 REST 端口连接器的创建、暂停、删除均通过该端口操作见下文 3.3。Converter 采用字节透传key.converter与value.converter均使用ByteArrayConverter即消息以原始字节进出由连接器端自行解析。这与 ClickHouse 连接器支持任意编码manifest 中supported_encodings: any的特性是配套的——具体解析格式由连接器配置如format参数决定。内部状态主题_connectors_offsets、_connectors_configs、_connectors_status三个主题用于 Kafka Connect 自身存储偏移量、配置与状态replication.factor-1表示跟随 Broker 默认副本数配置。3.2 ClickHouse 服务端配合连接器写入的是 ClickHouse 表因此 ClickHouse 服务的可达性至关重要。本仓库的 ClickHouse 配置extra/clickhouse/config.d/config.xml 声明了listen_host 0.0.0.0、TCP 端口9000、HTTP 端口8123即同时开放原生 TCP 协议连接器通常使用该协议与 HTTP 协议extra/clickhouse/users.d/users.xml 为default用户设置了密码default并开放::/0网段访问同时开启了access_management、named_collection_control等管理能力。连接器配置中需与之对应填写 ClickHouse 地址本仓库 compose 网络中主机名为clickhouse、端口9000、用户名与密码。3.3 注册 Sink 连接器REST API 方式Kafka Connect 运行后通过 REST API 注册连接器实例格式为 JSONcurl -X POST http://localhost:8083/connectors \ -H Content-Type: application/json \ -d { name: clickhouse-sink-events, config: { connector.class: com.clickhouse.kafka.connect.ClickHouseSinkConnector, tasks.max: 1, topics: events_raw, hostname: clickhouse, port: 9000, username: default, password: default, database: default, table: events, format: JSONEachRow, key.converter: org.apache.kafka.connect.converters.ByteArrayConverter, value.converter: org.apache.kafka.connect.converters.ByteArrayConverter } }提示以上connector.class、hostname、port、username、password、database、table、format等参数是 ClickHouse Sink 连接器常见配置项具体参数名与取值范围请以官方文档为准。本仓库 extra/debezium_config.json 展示了同构的完整连接器配置样例Debezium PostgreSQL CDC Source其中name、connector.class、tasks.max、converter、transforms 的组织方式与 Sink 连接器一致可作为在 Kafka Connect REST API 中提交 JSON 配置的格式参考。四、KeyToValue Transformation把 Kafka 消息 Key 写入 ClickHouse 列这是官方 README 花费篇幅最多的核心实战功能也是本文的重点。4.1 为什么需要它Kafka 消息由Key与Value两部分组成。默认情况下Sink 连接器写入 ClickHouse 的是消息 Value负载消息 Key 通常被丢弃。但在很多场景下Key 是有价值的业务字段——例如事件消息的 Key 可能是事件 ID、组织 ID 或会话 ID。KeyToValue这个官方提供的 Transformation 的作用正是把 Kafka 消息的 Key 作为一条记录转换进 Value 流中从而让 Key 能够以独立列的形式写入 ClickHouse。4.2 建表为目标列预留字段要在 ClickHouse 中保存 Key需要先在目标表里显式定义对应列。官方 README 给出的建表模板如下_key为默认列名类型为 StringCREATE TABLE your_table_name ( your_column_name String, ... ... ... _key String ) ENGINE MergeTree()要点说明列名与类型默认列名为_key类型为String。由于 Kafka 消息 Key 在字节层面是任意二进制将其声明为String是最稳妥的做法。引擎选择模板使用MergeTree()。在生产分析场景中可根据去重、排序等需求换用ReplacingMergeTree、AggregatingMergeTree等 MergeTree 家族引擎。务必保证连接器配置中的field与表列名一致否则写入会因列缺失而失败。4.3 连接器配置启用 Transformation在连接器配置中加入如下三行即可启用该转换transformskeyToValue transforms.keyToValue.typecom.clickhouse.kafka.connect.transforms.KeyToValue transforms.keyToValue.field_key参数逐行解读配置项含义transformskeyToValue声明启用名为keyToValue的 SMT 链若有多个转换用逗号分隔transforms.keyToValue.typecom.clickhouse.kafka.connect.transforms.KeyToValue指定 SMT 实现类全限定名transforms.keyToValue.field_key指定转换后的字段名须与目标表列名一致默认值即_key这套机制之所以可用正是得益于 manifest.json 中声明的single_message_transforms: true单条消息转换特性——SMT 会在 Sink 写入前对每条记录执行转换将 Key 注入记录体中。提示field参数可以改成任意自定义列名例如transforms.keyToValue.fieldmsg_key但务必同步修改建表 SQL 中的列名。五、性能测试官方 benchmark 工程官方 README 说明连接器仓库中附带一个独立的Gradle 工程benchmark专门用于该连接器的性能测试其运行方法与详细说明见该工程自身的 README。对使用者来说理解 benchmark 的意义在于基准验证在上线前用接近真实负载的数据量压测验证 Kafka 吞吐、ClickHouse 写入并发与连接器tasks.max的匹配度参数调优依据结合压测结果调整批次大小、缓冲配置与 ClickHouse 侧写入设置。说明benchmark 工程属于连接器官方仓库Lago 仓库随附的 clickhouse-kafka-connect-v1.3.4 目录只包含运行产物jar、manifest、README、LICENSE不包含 benchmark 源码如需运行基准测试请从连接器官方仓库获取。六、使用建议与常见问题6.1 生产配置建议版本匹配连接器要求 ClickHouse v22.5 及以上见 manifest.json本仓库使用的clickhouse-server:26.2-alpine满足要求Converter 一致性确保 Kafka Connect 全局 converter本仓库为ByteArrayConverter与连接器期望的输入格式匹配避免字节解码错误事务与 Exactly-Once开启 Exactly-Once 语义需要 Kafka Broker 端transaction.state.log等事务支持且会引入额外开销请结合业务对精确性的要求决定是否启用监控与状态Kafka Connect 的_connectors_offsets等状态主题、ClickHouse 侧的 monitoring.md 与日志配置extra/clickhouse/config.d/config.xml 中logger段可作为排查写入异常的依据。6.2 常见问题排查思路现象排查方向连接器任务反复失败检查 ClickHouse 可达性hostname/port、用户名密码、目标库表是否存在Key 列没有数据确认transformskeyToValue是否在配置中生效、field是否与表列名一致写入延迟高检查tasks.max并行度、ClickHouse 端写入批次与 Kafka 侧消费速率数据重复确认交付语义选择必要时在 ClickHouse 侧改用去重引擎并在读取时去重6.3 获取帮助官方 README 建议使用中遇到问题可在连接器官方仓库的 Issue 区提交问题或在 ClickHouse 官方 Slack 中提问官方维护团队会持续跟进。本仓库随附的 doc/LICENSEApache 2.0也说明该连接器可自由用于商业与非商业场景。七、总结ClickHouse Kafka Connect Sink 连接器v1.3.4为Kafka/Redpanda → ClickHouse的实时数据管道提供了官方、开箱即用的落地方案通过KeyToValueSMT 可以低门槛地把 Kafka 消息 Key 保存为 ClickHouse 独立列配合本仓库 docker-compose.dev.yml 中已经搭好的 Kafka Connect 与 ClickHouse 环境即可快速构建事件分析数据管道。对于 Lago 这样以事件为计量核心的计费平台这套连接器是打通消息队列 → 分析存储的关键一环——将 Kafka 侧的事件流可靠地沉淀到 ClickHouse为后续的用量聚合、计费分析提供数据底座。【免费下载链接】lagoOpen Source Metering and Usage Based Billing API ⭐️ Consumption tracking, Subscription management, Pricing iterations, Payment orchestration Revenue analytics项目地址: https://gitcode.com/GitHub_Trending/la/lago创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表