跳转到主要内容
Kafka 表引擎可用于从 Apache Kafka 和其他兼容 Kafka API 的消息代理 (例如 Redpanda、Amazon MSK) 读取数据,也可向其写入数据

Kafka 到 ClickHouse

如果您使用的是 ClickHouse Cloud,我们建议改用 ClickPipes。ClickPipes 原生支持私有网络连接,可对摄取和集群资源分别进行扩缩容,并为流式 Kafka 数据摄取到 ClickHouse 提供全面监控。
要使用 Kafka 表引擎,您应当对 ClickHouse materialized views 有较为全面的了解。

概述

首先,我们关注最常见的用例:使用 Kafka 表引擎 将数据从 Kafka 插入 ClickHouse。 Kafka 表引擎 允许 ClickHouse 直接从 Kafka topic 读取数据。虽然这对于查看某个 topic 上的消息很有用,但该引擎在设计上只支持一次性读取。也就是说,当对该表发出查询时,它会从队列中消费数据,并在将结果返回给调用方之前推进消费者偏移量。实际上,如果不重置这些偏移量,数据就无法再次读取。 要将通过 table engine 读取到的数据持久化,我们需要一种机制来捕获这些数据并将其插入到另一张表中。基于触发器的 materialized view 原生提供了这一能力。materialized view 会触发对 table engine 的读取,并接收成批的文档。TO 子句决定数据的目标端——通常是一张属于 MergeTree 家族 的表。如下图所示:

步骤

1. 准备
如果目标 topic 中已经有数据,你可以基于下面的内容调整后用于自己的数据集。或者,也可以使用这里提供的 GitHub 示例数据集。下面的示例使用的就是这个数据集;为简洁起见,相比这里提供的完整数据集,它采用了精简后的 schema 和部分行 (具体来说,我们只保留了与 ClickHouse 仓库 相关的 GitHub 事件) 。不过,这仍足以让随该数据集发布的大多数查询正常运行。
2. 配置 ClickHouse
如果要连接到启用安全机制的 Kafka,则必须执行此步骤。这些设置无法通过 SQL DDL 命令传入,必须在 ClickHouse 的 config.xml 中配置。这里假设你连接的是启用了 SASL 的 instance。与 Confluent Cloud 交互时,这是最简单的方法。
可以将上述代码片段放到 conf.d/ 目录下的新文件中,或合并到现有配置文件中。有关可配置的设置,请参见此处 我们还将创建一个名为 KafkaEngine 的数据库,用于本教程:
创建数据库后,你需要切换到该数据库:
3. 创建目标表
准备好目标表。为简洁起见,下面的示例使用了精简版的 GitHub schema。请注意,虽然这里使用的是 MergeTree 表引擎,但此示例也很容易改写为适用于 MergeTree 家族 中的任何成员。
4. 创建并填充 topic
接下来,我们将创建一个 topic。可以使用多种工具来完成这一步。如果 Kafka 运行在本地机器上,或在 Docker 容器中运行,RPK 是一个不错的选择。我们可以运行以下命令来创建一个名为 github、包含 5 个分区的 topic:
如果我们使用的是 Confluent Cloud 上的 Kafka,可能会更倾向于使用 Confluent CLI
现在我们需要向这个 topic 写入一些数据,这里将使用 kcat。如果是在本地运行 Kafka 且未启用身份验证,可以运行类似下面的命令:
或者,如果 Kafka 集群使用 SASL 进行身份验证,请使用以下内容:
该数据集包含 200,000 行,因此几秒内即可完成摄取。如果你想处理更大的数据集,请参阅 ClickHouse/kafka-samples GitHub 代码仓库中的大数据集部分
5. 创建 Kafka 表引擎
下面的示例创建了一个表引擎,其 schema 与 merge tree 表相同。这并非绝对必要,因为你可以在目标表中使用别名或临时列。不过,这些设置很重要——请注意,这里使用 JSONEachRow 作为从 Kafka topic 消费 JSON 时的数据类型。githubclickhouse 这两个值分别表示 topic 名称和消费者组名称。实际上,topics 也可以是一个值列表。
我们将在下文讨论引擎设置和性能调优。此时,对表 github_queue 执行一个简单的 select 查询后,应该能读取到一些行。请注意,这会将消费者偏移量向前移动,因此如果不进行重置,就无法再次读取这些行。另请注意其中的 limit,以及必需参数 stream_like_engine_allow_direct_select.
6. 创建 materialized view
materialized view 会连接前面创建的两个表:从 Kafka table engine 读取数据,并将其插入目标 merge tree 表。我们可以在此过程中进行多种数据转换。这里我们只做一次简单的读取和插入。使用 * 的前提是列名完全一致 (区分大小写) 。
创建后,materialized view 会连接到 Kafka 引擎并开始读取,将行插入目标表。该过程会持续运行,后续插入到 Kafka 的消息也会被持续消费。你可以根据需要重新运行插入脚本,向 Kafka 再插入更多消息。
7. 确认数据行已插入
确认目标表中已有数据:
你应该会看到 200,000 行数据:

常用操作

停止和恢复消息消费
要停止消息消费,可以分离 Kafka 引擎表:
这不会影响消费者组的偏移量。要重启消费并从之前的偏移量继续,请重新附加该表。
添加 Kafka 元数据
将原始 Kafka 消息中的元数据在摄取到 ClickHouse 后保留下来,通常会很有帮助。例如,我们可能希望了解某个特定 topic 或分区已经消费了多少。为此,Kafka 表引擎提供了多个虚拟列。通过修改 schema 和 materialized view 的 SELECT 语句,可以将这些虚拟列作为普通列持久化到目标表中。 首先,在向目标表添加列之前,先执行上文所述的停止操作。
下面我们添加信息列,用于标识来源 topic 以及该行来自哪个分区。
接下来,我们需要确保虚拟列已按要求映射。 虚拟列带有 _ 前缀。 虚拟列的完整列表可在此处查看。 要使用这些虚拟列更新表,我们需要删除 materialized view,重新 Attach Kafka 引擎表,并重新创建 materialized view。
新读取的行应包含这些元数据。
结果如下:
修改 Kafka 引擎设置
我们建议删除 Kafka 引擎表,并使用新设置重新创建。在此过程中,无需修改 materialized view——Kafka 引擎表重建后,消息消费会自动恢复。
调试问题
身份验证等错误不会出现在 Kafka 引擎 DDL 的响应中。要诊断此类问题,建议查看 ClickHouse 的主日志文件 clickhouse-server.err.log。还可以通过配置为底层 Kafka 客户端库 librdkafka 启用更详细的 trace 日志。
处理格式错误的消息
Kafka 常常被当作数据“堆放场”使用。这会导致 topic 中混杂着不同的消息格式和不一致的字段名。应尽量避免这种情况,并利用 Kafka 的功能 (例如 Kafka Streams 或 ksqlDB) ,确保消息在写入 Kafka 之前就是格式良好且一致的。如果无法采用这些方案,ClickHouse 也提供了一些可用于缓解问题的功能。
  • 将消息字段按字符串处理。如有需要,可以在 materialized view 语句中使用函数进行清洗和类型转换。虽然这不应视为生产环境方案,但对于一次性摄取可能会有帮助。
  • 如果你从某个 topic 中消费 JSON,并使用 JSONEachRow format,请使用设置 input_format_skip_unknown_fields。写入数据时,默认情况下,如果输入数据包含目标表中不存在的列,ClickHouse 会抛出异常。但如果启用此选项,这些多出的列会被忽略。同样,这也不是生产级方案,而且可能会让其他人感到困惑。
  • 可以考虑使用设置 kafka_skip_broken_messages。该设置要求用户为每个块中格式错误的消息指定容忍度,并结合 kafka_max_block_size 来判断。如果超过这个容忍度 (按消息绝对数量计算) ,则会恢复默认的异常行为,并跳过其他消息。
投递语义以及重复数据带来的挑战
Kafka 表引擎具有至少一次 (at-least-once) 投递语义。在一些已知但罕见的情况下,可能会出现重复数据。例如,消息可能已经从 Kafka 读取并成功插入 ClickHouse。但在提交新的 偏移量 之前,与 Kafka 的连接丢失了。在这种情况下,就需要重试该块。如果将分布式表或 ReplicatedMergeTree 用作目标表,则该块可以去重。虽然这会降低重复行出现的概率,但它依赖于块完全一致。像 Kafka 再均衡这样的事件可能会破坏这一前提,从而在少数情况下导致重复数据。
基于仲裁的插入
在 ClickHouse 中,如果需要更高的投递保障,可能需要使用基于仲裁的插入。这项设置不能在 materialized view 或目标表上配置,但可以为用户 profile 设置,例如:

ClickHouse 到 Kafka

虽然这种用例较为少见,但也可以将 ClickHouse 数据持久化到 Kafka 中。例如,我们将手动向 Kafka 表引擎插入行。随后,同一个 Kafka 引擎会读取这些数据,其 materialized view 会将数据写入 MergeTree 表。最后,我们将演示在向 Kafka 插入数据时如何使用 materialized views,从现有 source table 中读取数据。

步骤

我们的初始目标如下图所示: 我们假设你已按照 Kafka to ClickHouse 中的步骤创建好这些表和视图,并且该 topic 中的数据已被完全消费。
1. 直接插入行
首先,确认目标表中的行数。
此时应有 200,000 行:
现在,将 GitHub 目标表中的行重新插入到 Kafka 表引擎 github_queue 中。请注意,这里使用的是 JSONEachRow 格式,并将 SELECT 的 LIMIT 设为 100。
重新统计 GitHub 表中的行数,确认其已增加 100。如上图所示,数据行先通过 Kafka 表引擎写入 Kafka,随后再由同一引擎重新读取,并由我们的 materialized view 插入到 GitHub 目标表中!
你应该会看到新增了 100 行:
2. 使用 materialized view
当文档插入表中时,我们可以利用 materialized view 将消息推送到 Kafka 引擎 (以及某个 topic) 。当行插入 GitHub 表时,会触发一个 materialized view,进而将这些行重新插入到 Kafka 引擎中,并写入一个新的 topic。如下图所示: 创建一个新的 Kafka topic github_out 或等效项。确保 Kafka 表引擎 github_out_queue 指向该 topic。
现在创建一个新的 materialized view github_out_mv,使其指向 GitHub 表,并在触发时将行插入到上述引擎中。这样一来,添加到 GitHub 表中的内容就会被推送到新的 Kafka topic。
如果你向原始的 github topic (在 Kafka to ClickHouse 中创建) 插入数据,文档就会自动出现在 “github_clickhouse” topic 中。你可以使用原生 Kafka 工具来确认这一点。例如,下面我们使用 kcat 向由 Confluent Cloud 托管的 github topic 插入 100 行数据:
读取 github_out topic 应可确认消息已送达。
虽然这是一个较复杂的示例,但它展示了 materialized view 与 Kafka 引擎结合使用时的强大能力。

集群与性能

使用 ClickHouse 集群

通过 Kafka 消费者组,多个 ClickHouse 实例可以同时从同一个 topic 读取数据。每个消费者都会以 1:1 的映射关系分配到一个 topic 分区。在使用 Kafka 表引擎对 ClickHouse 的消费能力进行扩缩容时,请注意,集群中的消费者总数不能超过该 topic 的分区数。因此,请务必提前为 topic 配置好合适的分区方案。 多个 ClickHouse 实例也可以配置为使用同一个消费者组 id 从某个 topic 读取数据——该 id 在创建 Kafka 表引擎时指定。因此,每个实例都会从一个或多个分区读取数据,并将数据分段插入其本地目标表。目标表则可以进一步配置为使用 ReplicatedMergeTree 来处理数据重复。这种方法可以让 Kafka 读取能力随着 ClickHouse 集群一同扩展,前提是 Kafka 有足够多的分区。

性能调优

在尝试提升 Kafka Engine 表的吞吐性能时,请考虑以下几点:
  • 性能会因消息大小、格式以及目标表类型而异。对于单个表引擎,达到 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 消费者滞后,以确保扩容得当。

其他设置

除了上文介绍的设置外,以下内容也值得关注:
  • Kafka_max_wait_ms - 重试前从 Kafka 读取消息的等待时间,以毫秒为单位。在用户 profile 级别设置,默认值为 5000。
底层 librdkafka 的所有设置 也可以放在 ClickHouse 配置文件中的 kafka 元素内——设置名称应写成 XML 元素,并将句点替换为下划线,例如:
这些是专家级设置,建议参考 Kafka 文档了解更深入的说明。
最后修改于 2026年6月12日