Kafka 到 ClickHouse
如果您使用的是 ClickHouse Cloud,我们建议改用 ClickPipes。ClickPipes 原生支持私有网络连接,可对摄取和集群资源分别进行扩缩容,并为流式 Kafka 数据摄取到 ClickHouse 提供全面监控。
概述
TO 子句决定数据的目标端——通常是一张属于 MergeTree 家族 的表。如下图所示:
步骤
1
准备
如果你的目标 topic 中已有数据,你可以据此调整以下内容,以适配你的数据集。或者,你也可以使用这里提供的 GitHub 示例数据集。为简洁起见,下面的示例使用的就是这个数据集;与这里提供的完整数据集相比,它采用了精简版 schema,并且只包含部分行 (具体来说,我们仅保留与 ClickHouse 软件源 相关的 GitHub 事件) 。不过,这仍足以让随数据集发布的大多数查询正常运行。
2
配置 ClickHouse
如果您要连接到启用了安全机制的 Kafka,则此步骤必不可少。这些设置无法通过 SQL DDL 命令传入,必须在 ClickHouse 的 config.xml 中配置。这里假设您连接的是启用了 SASL 保护的实例。这是在与 Confluent Cloud 交互时最简单的方法。将上述代码片段放入 创建好数据库后,您需要切换到该数据库:
conf.d/ 目录下的新文件中,或者将其合并到现有配置文件中。有关可配置的设置,请参见此处。我们还将创建一个名为 KafkaEngine 的数据库,供本教程使用:3
创建目标表
准备好目标表。为简洁起见,下面的示例使用了精简版的 GitHub schema。请注意,虽然这里使用的是 MergeTree 表引擎,但此示例也很容易改写为适用于 MergeTree 家族 中的任何成员。
4
创建 topic 并写入数据
接下来,我们将创建一个 topic。为此可以使用多种工具。如果是在本机上本地运行 Kafka,或者在 Docker 容器中运行 Kafka,RPK 就很合适。我们可以运行以下命令来创建一个名为 如果我们是在 Confluent Cloud 上运行 Kafka,则可能更适合使用 Confluent 命令行客户端:现在我们需要向这个 topic 写入一些数据,这里将使用 kcat。如果你在本地运行 Kafka 且禁用了身份验证,可以运行类似下面的命令:或者,如果我们的 Kafka 集群使用 SASL 进行身份验证,则使用以下内容:该数据集包含 200,000 行,因此只需几秒钟即可完成摄取。如果你想使用更大的数据集,请查看 GitHub 仓库 ClickHouse/kafka-samples 中的大数据集部分。
github、具有 5 个分区的 topic:5
创建 Kafka 表引擎
下面的示例创建了一个与 MergeTree 表具有相同 schema 的表引擎。这并非严格必需,因为你可以在目标表中使用别名列或临时列。不过,这些设置很重要——请注意,这里使用 下面将讨论引擎设置和性能调优。此时,对表
JSONEachRow 作为从 Kafka topic 中消费 JSON 的数据类型。其中,github 和 clickhouse 分别表示 topic 名称和消费者组名称。实际上,topics 也可以是一个值列表。github_queue 执行一个简单的 select 应该能读出一些行。请注意,这会将消费者 offsets 向前推进,因此如果不进行重置,这些行将无法再次读取。另请注意限制以及必需参数 stream_like_engine_allow_direct_select.6
创建 materialized view
materialized view 会连接前面创建的两个表,从 Kafka 表引擎读取数据,并将其插入目标 MergeTree 表。我们可以进行多种数据转换。这里我们将执行简单的读取和插入操作。使用 * 的前提是列名完全一致 (区分大小写) 。在创建时,materialized view 会连接到 Kafka 引擎并开始读取,将行插入目标表。此过程会无限期持续,之后插入到 Kafka 的消息也会被消费。你可以随时重新运行插入脚本,向 Kafka 再插入更多消息。
7
确认行已插入
确认目标表中有数据:你应该能看到 200,000 行:
常用操作
停止和恢复消息消费
添加 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
直接插入数据行
首先,确认目标表中的行数。此时应有 200,000 行:现在将行从 GitHub 目标表重新插入到 Kafka 表引擎 github_queue 中。请注意,我们使用了 JSONEachRow 格式,并通过 LIMIT 将 select 限制为 100。重新统计 GitHub 表中的行数,以确认其已增加 100。正如上图所示,行先通过 Kafka 表引擎插入到 Kafka 中,随后再由同一个引擎重新读取,并由我们的 materialized view 插入到 GitHub 目标表中!你应该会看到另外 100 行:
2
使用 materialized view
当文档插入表中时,我们可以利用 materialized view 将消息推送到 Kafka 引擎 (以及某个 topic) 。当行插入 GitHub 表时,会触发一个 materialized view,进而将这些行重新插入到 Kafka 引擎中,并写入一个新的 topic。如下图所示:创建一个新的 Kafka topic 现在创建一个新的 materialized view 如果你向原始的 github topic (在 Kafka 到 ClickHouse 中创建) 插入数据,文档就会自动出现在 “github_clickhouse” topic 中。你可以使用原生 Kafka 工具来确认这一点。例如,下面我们使用 kcat 向由 Confluent Cloud 托管的 github topic 插入 100 行数据:对 尽管这是一个较为复杂的示例,但它充分展示了 materialized view 与 Kafka 引擎结合使用时的强大能力。
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。