Kafka 到 ClickHouse
如果您使用的是 ClickHouse Cloud,我们建议改用 ClickPipes。ClickPipes 原生支持私有网络连接,可对摄取和集群资源分别进行扩缩容,并为流式 Kafka 数据摄取到 ClickHouse 提供全面监控。
概述
TO 子句决定数据的目标端——通常是一张属于 MergeTree 家族 的表。如下图所示:
步骤
1. 准备
2. 配置 ClickHouse
conf.d/ 目录下的新文件中,或合并到现有配置文件中。有关可配置的设置,请参见此处。
我们还将创建一个名为 KafkaEngine 的数据库,用于本教程:
3. 创建目标表
4. 创建并填充 topic
github、包含 5 个分区的 topic:
5. 创建 Kafka 表引擎
JSONEachRow 作为从 Kafka topic 消费 JSON 时的数据类型。github 和 clickhouse 这两个值分别表示 topic 名称和消费者组名称。实际上,topics 也可以是一个值列表。
github_queue 执行一个简单的 select 查询后,应该能读取到一些行。请注意,这会将消费者偏移量向前移动,因此如果不进行重置,就无法再次读取这些行。另请注意其中的 limit,以及必需参数 stream_like_engine_allow_direct_select.
6. 创建 materialized view
7. 确认数据行已插入
常用操作
停止和恢复消息消费
添加 Kafka 元数据
_ 前缀。
虚拟列的完整列表可在此处查看。
要使用这些虚拟列更新表,我们需要删除 materialized view,重新 Attach Kafka 引擎表,并重新创建 materialized view。
修改 Kafka 引擎设置
调试问题
处理格式错误的消息
- 将消息字段按字符串处理。如有需要,可以在 materialized view 语句中使用函数进行清洗和类型转换。虽然这不应视为生产环境方案,但对于一次性摄取可能会有帮助。
- 如果你从某个 topic 中消费 JSON,并使用 JSONEachRow format,请使用设置
input_format_skip_unknown_fields。写入数据时,默认情况下,如果输入数据包含目标表中不存在的列,ClickHouse 会抛出异常。但如果启用此选项,这些多出的列会被忽略。同样,这也不是生产级方案,而且可能会让其他人感到困惑。 - 可以考虑使用设置
kafka_skip_broken_messages。该设置要求用户为每个块中格式错误的消息指定容忍度,并结合kafka_max_block_size来判断。如果超过这个容忍度 (按消息绝对数量计算) ,则会恢复默认的异常行为,并跳过其他消息。
投递语义以及重复数据带来的挑战
基于仲裁的插入
ClickHouse 到 Kafka
步骤
1. 直接插入行
2. 使用 materialized view
github_out 或等效项。确保 Kafka 表引擎 github_out_queue 指向该 topic。
github_out_mv,使其指向 GitHub 表,并在触发时将行插入到上述引擎中。这样一来,添加到 GitHub 表中的内容就会被推送到新的 Kafka topic。
github_out topic 应可确认消息已送达。
集群与性能
使用 ClickHouse 集群
性能调优
- 性能会因消息大小、格式以及目标表类型而异。对于单个表引擎,达到 100k 行/秒通常是可实现的。默认情况下,消息会按块读取,由参数
kafka_max_block_size控制。其默认值为 max_insert_block_size,默认为 1,048,576。除非消息特别大,否则几乎总是应该增大该值。500k 到 1M 的取值并不少见。请测试并评估其对吞吐性能的影响。 - 可以使用
kafka_num_consumers增加表引擎的消费者数量。不过,默认情况下,除非将kafka_thread_per_consumer从默认值 1 改为其他值,否则插入会被串行化到单个线程中。将其设为 1 以确保 flush 操作并行执行。请注意,创建一个具有 N 个消费者 (且kafka_thread_per_consumer=1) 的 Kafka 引擎表,在逻辑上等同于创建 N 个 Kafka 引擎,每个引擎各自配有一个 materialized view,且kafka_thread_per_consumer=0。 - 增加消费者并非没有代价。每个消费者都会维护自己的缓冲区和线程,从而增加 server 开销。如有可能,请先优先通过集群线性扩展来分摊负载,同时留意消费者带来的额外开销。
- 如果 Kafka 消息吞吐量波动较大且可以接受一定延迟,可考虑增大
stream_flush_interval_ms,以确保刷出更大的块。 - background_message_broker_schedule_pool_size 用于设置执行后台任务的线程数。这些线程会用于 Kafka 流式处理。该设置会在 ClickHouse server 启动时生效,且不能在用户 session 中更改,默认值为 16。如果你在日志中看到超时,适当增大该值可能是合适的。
- 与 Kafka 通信时使用的是
librdkafka库,而它本身也会创建线程。因此,大量 Kafka 表或消费者可能会导致大量上下文切换。可以将这部分负载分散到整个集群中,并尽可能只复制目标表;或者考虑使用一个表引擎从多个 topic 读取数据——支持值列表。单个表也可以被多个 materialized view 读取,每个视图分别过滤特定 topic 的数据。
其他设置
- Kafka_max_wait_ms - 重试前从 Kafka 读取消息的等待时间,以毫秒为单位。在用户 profile 级别设置,默认值为 5000。