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

步骤
准备
如果你的目标 topic 中已有数据,你可以据此调整以下内容,以适配你的数据集。或者,你也可以使用这里提供的 GitHub 示例数据集。为简洁起见,下面的示例使用的就是这个数据集;与这里提供的完整数据集相比,它采用了精简版 schema,并且只包含部分行 (具体来说,我们仅保留与 ClickHouse 软件源 相关的 GitHub 事件) 。不过,这仍足以让随数据集发布的大多数查询正常运行。
配置 ClickHouse
如果您要连接到启用了安全机制的 Kafka,则此步骤必不可少。这些设置无法通过 SQL DDL 命令传入,必须在 ClickHouse 的 config.xml 中配置。这里假设您连接的是启用了 SASL 保护的实例。这是在与 Confluent Cloud 交互时最简单的方法。
<clickhouse>
<kafka>
<sasl_username>username</sasl_username>
<sasl_password>password</sasl_password>
<security_protocol>sasl_ssl</security_protocol>
<sasl_mechanisms>PLAIN</sasl_mechanisms>
</kafka>
</clickhouse>将上述代码片段放入 conf.d/ 目录下的新文件中,或者将其合并到现有配置文件中。有关可配置的设置,请参见此处。
我们还将创建一个名为 KafkaEngine 的数据库,供本教程使用:
CREATE DATABASE KafkaEngine;创建好数据库后,您需要切换到该数据库:
USE KafkaEngine;创建目标表
准备好目标表。为简洁起见,下面的示例使用了精简版的 GitHub schema。请注意,虽然这里使用的是 MergeTree 表引擎,但此示例也很容易改写为适用于 MergeTree 家族 中的任何成员。
CREATE TABLE github
(
file_time DateTime,
event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
actor_login LowCardinality(String),
repo_name LowCardinality(String),
created_at DateTime,
updated_at DateTime,
action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
comment_id UInt64,
path String,
ref LowCardinality(String),
ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
creator_user_login LowCardinality(String),
number UInt32,
title String,
labels Array(LowCardinality(String)),
state Enum('none' = 0, 'open' = 1, 'closed' = 2),
assignee LowCardinality(String),
assignees Array(LowCardinality(String)),
closed_at DateTime,
merged_at DateTime,
merge_commit_sha String,
requested_reviewers Array(LowCardinality(String)),
merged_by LowCardinality(String),
review_comments UInt32,
member_login LowCardinality(String)
) ENGINE = MergeTree ORDER BY (event_type, repo_name, created_at)创建 topic 并写入数据
接下来,我们将创建一个 topic。为此可以使用多种工具。如果是在本机上本地运行 Kafka,或者在 Docker 容器中运行 Kafka,RPK 就很合适。我们可以运行以下命令来创建一个名为 github、具有 5 个分区的 topic:
rpk topic create -p 5 github --brokers <host>:<port>如果我们是在 Confluent Cloud 上运行 Kafka,则可能更适合使用 Confluent 命令行客户端:
confluent kafka topic create --if-not-exists github现在我们需要向这个 topic 写入一些数据,这里将使用 kcat。如果你在本地运行 Kafka 且禁用了身份验证,可以运行类似下面的命令:
cat github_all_columns.ndjson |
kcat -P \
-b <host>:<port> \
-t github或者,如果我们的 Kafka 集群使用 SASL 进行身份验证,则使用以下内容:
cat github_all_columns.ndjson |
kcat -P \
-b <host>:<port> \
-t github
-X security.protocol=sasl_ssl \
-X sasl.mechanisms=PLAIN \
-X sasl.username=<username> \
-X sasl.password=<password> \该数据集包含 200,000 行,因此只需几秒钟即可完成摄取。如果你想使用更大的数据集,请查看 GitHub 仓库 ClickHouse/kafka-samples 中的大数据集部分。
创建 Kafka 表引擎
下面的示例创建了一个与 MergeTree 表具有相同 schema 的表引擎。这并非严格必需,因为你可以在目标表中使用别名列或临时列。不过,这些设置很重要——请注意,这里使用 JSONEachRow 作为从 Kafka topic 中消费 JSON 的数据类型。其中,github 和 clickhouse 分别表示 topic 名称和消费者组名称。实际上,topics 也可以是一个值列表。
CREATE TABLE github_queue
(
file_time DateTime,
event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
actor_login LowCardinality(String),
repo_name LowCardinality(String),
created_at DateTime,
updated_at DateTime,
action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
comment_id UInt64,
path String,
ref LowCardinality(String),
ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
creator_user_login LowCardinality(String),
number UInt32,
title String,
labels Array(LowCardinality(String)),
state Enum('none' = 0, 'open' = 1, 'closed' = 2),
assignee LowCardinality(String),
assignees Array(LowCardinality(String)),
closed_at DateTime,
merged_at DateTime,
merge_commit_sha String,
requested_reviewers Array(LowCardinality(String)),
merged_by LowCardinality(String),
review_comments UInt32,
member_login LowCardinality(String)
)
ENGINE = Kafka('kafka_host:9092', 'github', 'clickhouse',
'JSONEachRow') SETTINGS kafka_thread_per_consumer = 0, kafka_num_consumers = 1;下面将讨论引擎设置和性能调优。此时,对表 github_queue 执行一个简单的 select 应该能读出一些行。请注意,这会将消费者 offsets 向前推进,因此如果不进行重置,这些行将无法再次读取。另请注意限制以及必需参数 stream_like_engine_allow_direct_select.
创建 materialized view
materialized view 会连接前面创建的两个表,从 Kafka 表引擎读取数据,并将其插入目标 MergeTree 表。我们可以进行多种数据转换。这里我们将执行简单的读取和插入操作。使用 * 的前提是列名完全一致 (区分大小写) 。
CREATE MATERIALIZED VIEW github_mv TO github AS
SELECT *
FROM github_queue;在创建时,materialized view 会连接到 Kafka 引擎并开始读取,将行插入目标表。此过程会无限期持续,之后插入到 Kafka 的消息也会被消费。你可以随时重新运行插入脚本,向 Kafka 再插入更多消息。
确认行已插入
确认目标表中有数据:
SELECT count() FROM github;你应该能看到 200,000 行:
┌─count()─┐
│ 200000 │
└─────────┘常用操作
停止和恢复消息消费
要停止消息消费,可以分离 Kafka 引擎表:
DETACH TABLE github_queue;这不会影响消费者组的偏移量。要重启消费并从之前的偏移量继续,请重新附加该表。
ATTACH TABLE github_queue;添加 Kafka 元数据
将原始 Kafka 消息中的元数据在摄取到 ClickHouse 后保留下来,通常会很有帮助。例如,我们可能希望了解某个特定 topic 或分区已经消费了多少。为此,Kafka 表引擎提供了多个虚拟列。通过修改 schema 和 materialized view 的 select 语句,可以将这些虚拟列作为普通列持久化到目标表中。
首先,在向目标表添加列之前,先执行上文所述的停止操作。
DETACH TABLE github_queue;下面我们添加信息列,用于标识来源 topic 以及该行来自哪个分区。
ALTER TABLE github
ADD COLUMN topic String,
ADD COLUMN partition UInt64;接下来,我们需要确保虚拟列已按要求映射。
虚拟列带有 _ 前缀。
虚拟列的完整列表可在此处查看。
要使用这些虚拟列更新表,我们需要删除 materialized view,重新 Attach Kafka 引擎表,并重新创建 materialized view。
DROP VIEW github_mv;ATTACH TABLE github_queue;CREATE MATERIALIZED VIEW github_mv TO github AS
SELECT *, _topic AS topic, _partition as partition
FROM github_queue;新读取的行应包含这些元数据。
SELECT actor_login, event_type, created_at, topic, partition
FROM github
LIMIT 10;结果如下:
| actor_login | event_type | created_at | topic | partition |
|---|---|---|---|---|
| IgorMinar | CommitCommentEvent | 2011-02-12 02:22:00 | github | 0 |
| queeup | CommitCommentEvent | 2011-02-12 02:23:23 | github | 0 |
| IgorMinar | CommitCommentEvent | 2011-02-12 02:23:24 | github | 0 |
| IgorMinar | CommitCommentEvent | 2011-02-12 02:24:50 | github | 0 |
| IgorMinar | CommitCommentEvent | 2011-02-12 02:25:20 | github | 0 |
| dapi | CommitCommentEvent | 2011-02-12 06:18:36 | github | 0 |
| sourcerebels | CommitCommentEvent | 2011-02-12 06:34:10 | github | 0 |
| jamierumbelow | CommitCommentEvent | 2011-02-12 12:21:40 | github | 0 |
| jpn | CommitCommentEvent | 2011-02-12 12:24:31 | github | 0 |
| Oxonium | CommitCommentEvent | 2011-02-12 12:31:28 | github | 0 |
修改 Kafka 引擎设置
我们建议删除 Kafka 引擎表,并使用新设置重新创建。在此过程中,无需修改 materialized view——Kafka 引擎表重建后,消息消费会自动恢复。
调试问题
身份验证等错误不会出现在 Kafka 引擎 DDL 的响应中。要诊断此类问题,建议查看 ClickHouse 的主日志文件 clickhouse-server.err.log。还可以通过配置为底层 Kafka 客户端库 librdkafka 启用更详细的 trace 日志。
<kafka>
<debug>all</debug>
</kafka>处理格式错误的消息
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 设置,例如:
<profiles>
<default>
<insert_quorum>2</insert_quorum>
</default>
</profiles>ClickHouse 到 Kafka
虽然这种用例较为少见,但也可以将 ClickHouse 数据持久化到 Kafka 中。例如,我们将手动向 Kafka 表引擎插入行。随后,同一个 Kafka 引擎会读取这些数据,其 materialized view 会将数据写入 MergeTree 表。最后,我们将演示在向 Kafka 插入数据时如何使用 materialized views,从现有 source table 中读取数据。
步骤
我们的初始目标如下图所示:

我们假设你已按照 Kafka 到 ClickHouse 中的步骤创建好这些表和视图,并且该 topic 中的数据已被完全消费。
直接插入数据行
首先,确认目标表中的行数。
SELECT count() FROM github;此时应有 200,000 行:
┌─count()─┐
│ 200000 │
└─────────┘现在将行从 GitHub 目标表重新插入到 Kafka 表引擎 github_queue 中。请注意,我们使用了 JSONEachRow 格式,并通过 LIMIT 将 select 限制为 100。
INSERT INTO github_queue SELECT * FROM github LIMIT 100 FORMAT JSONEachRow重新统计 GitHub 表中的行数,以确认其已增加 100。正如上图所示,行先通过 Kafka 表引擎插入到 Kafka 中,随后再由同一个引擎重新读取,并由我们的 materialized view 插入到 GitHub 目标表中!
SELECT count() FROM github;你应该会看到另外 100 行:
┌─count()─┐
│ 200100 │
└─────────┘使用 materialized view
当文档插入表中时,我们可以利用 materialized view 将消息推送到 Kafka 引擎 (以及某个 topic) 。当行插入 GitHub 表时,会触发一个 materialized view,进而将这些行重新插入到 Kafka 引擎中,并写入一个新的 topic。如下图所示:

创建一个新的 Kafka topic github_out 或等效项。确保 Kafka 表引擎 github_out_queue 指向该 topic。
CREATE TABLE github_out_queue
(
file_time DateTime,
event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
actor_login LowCardinality(String),
repo_name LowCardinality(String),
created_at DateTime,
updated_at DateTime,
action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
comment_id UInt64,
path String,
ref LowCardinality(String),
ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
creator_user_login LowCardinality(String),
number UInt32,
title String,
labels Array(LowCardinality(String)),
state Enum('none' = 0, 'open' = 1, 'closed' = 2),
assignee LowCardinality(String),
assignees Array(LowCardinality(String)),
closed_at DateTime,
merged_at DateTime,
merge_commit_sha String,
requested_reviewers Array(LowCardinality(String)),
merged_by LowCardinality(String),
review_comments UInt32,
member_login LowCardinality(String)
)
ENGINE = Kafka('host:port', 'github_out', 'clickhouse_out',
'JSONEachRow') SETTINGS kafka_thread_per_consumer = 0, kafka_num_consumers = 1;现在创建一个新的 materialized view github_out_mv,使其指向 GitHub 表,并在触发时将行插入到上述引擎中。这样一来,添加到 GitHub 表中的内容就会被推送到新的 Kafka topic。
CREATE MATERIALIZED VIEW github_out_mv TO github_out_queue AS
SELECT file_time, event_type, actor_login, repo_name,
created_at, updated_at, action, comment_id, path,
ref, ref_type, creator_user_login, number, title,
labels, state, assignee, assignees, closed_at, merged_at,
merge_commit_sha, requested_reviewers, merged_by,
review_comments, member_login
FROM github
FORMAT JsonEachRow;如果你向原始的 github topic (在 Kafka 到 ClickHouse 中创建) 插入数据,文档就会自动出现在 "github_clickhouse" topic 中。你可以使用原生 Kafka 工具来确认这一点。例如,下面我们使用 kcat 向由 Confluent Cloud 托管的 github topic 插入 100 行数据:
head -n 10 github_all_columns.ndjson |
kcat -P \对 github_out topic 执行读取操作,即可确认消息已成功投递。
kcat -C \尽管这是一个较为复杂的示例,但它充分展示了 materialized view 与 Kafka 引擎结合使用时的强大能力。
集群与性能
使用 ClickHouse 集群
通过 Kafka 消费者组,多个 ClickHouse 实例可以同时从同一个 topic 读取数据。每个消费者都会以 1:1 的映射关系分配到一个 topic 分区。在使用 Kafka 表引擎对 ClickHouse 的消费能力进行扩缩容时,请注意,集群中的消费者总数不能超过该 topic 的分区数。因此,请务必提前为 topic 配置好合适的分区方案。
多个 ClickHouse 实例也可以配置为使用同一个消费者组 id 从某个 topic 读取数据——该 id 在创建 Kafka 表引擎时指定。因此,每个实例都会从一个或多个分区读取数据,并将数据分段插入其本地目标表。目标表则可以进一步配置为使用 ReplicatedMergeTree 来处理数据重复。这种方法可以让 Kafka 读取能力随着 ClickHouse 集群一同扩展,前提是 Kafka 有足够多的分区。

性能调优
在尝试提升 Kafka 引擎表的吞吐性能时,请考虑以下几点:
- 性能会因消息大小、格式以及目标表类型而异。对于单个表引擎,达到 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 元素,并将句点替换为下划线,例如:
<clickhouse>
<kafka>
<enable_ssl_certificate_verification>false</enable_ssl_certificate_verification>
</kafka>
</clickhouse>这些是专家级设置,建议参考 Kafka 文档了解更深入的说明。