> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-detect-table-modification.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

> 使用 Kafka 表引擎

# 使用 Kafka 表引擎

export const Image = ({img, alt, size = "lg"}) => {
  const normalizedSize = ["sm", "md", "lg"].includes(size) ? size : "lg";
  return <div className={`ch-image-${normalizedSize}`}>
      <Frame>
        <img src={img} alt={alt} />
      </Frame>
    </div>;
};

Kafka 表引擎可用于从 Apache Kafka 和其他兼容 Kafka API 的消息代理 (例如 Redpanda、Amazon MSK) [**读取**数据](#kafka-to-clickhouse)，也可向其[**写入**数据](#clickhouse-to-kafka)。

<div id="kafka-to-clickhouse">
  ### Kafka 到 ClickHouse
</div>

<Note>
  如果您使用的是 ClickHouse Cloud，我们建议改用 [ClickPipes](/zh/integrations/clickpipes/home)。ClickPipes 原生支持私有网络连接，可对摄取和集群资源分别进行扩缩容，并为流式 Kafka 数据摄取到 ClickHouse 提供全面监控。
</Note>

要使用 Kafka 表引擎，您应当对 [ClickHouse materialized views](/zh/concepts/features/materialized-views/cascading-materialized-views) 有较为全面的了解。

<div id="overview">
  #### 概述
</div>

首先，我们关注最常见的用例：使用 Kafka 表引擎 将数据从 Kafka 插入 ClickHouse。

Kafka 表引擎 允许 ClickHouse 直接从 Kafka topic 读取数据。虽然这对于查看某个 topic 上的消息很有用，但该引擎在设计上只支持一次性读取。也就是说，当对该表发出查询时，它会从队列中消费数据，并在将结果返回给调用方之前推进消费者偏移量。实际上，如果不重置这些偏移量，数据就无法再次读取。

要将通过 表引擎 读取到的数据持久化，我们需要一种机制来捕获这些数据并将其插入到另一张表中。基于触发器的 materialized view 原生提供了这一能力。materialized view 会触发对 表引擎 的读取，并接收成批的文档。`TO` 子句决定数据的目标端——通常是一张属于 [MergeTree 家族](/zh/reference/engines/table-engines/mergetree-family/index) 的表。如下图所示：

<Image img="https://mintcdn.com/private-7c7dfe99-detect-table-modification/j2pAbv7ihJZXp9qi/images/integrations/data-ingestion/kafka/kafka_01.webp?fit=max&auto=format&n=j2pAbv7ihJZXp9qi&q=85&s=d0459f397e0f1770cc72f32787693008" size="lg" alt="Kafka 表引擎 架构图" style={{width: '80%'}} width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_01.webp" />

<div id="steps">
  #### 步骤
</div>

<div id="1-prepare">
  ##### 1. 准备
</div>

如果目标 topic 中已经有数据，你可以基于下面的内容调整后用于自己的数据集。或者，也可以使用[这里](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson)提供的 GitHub 示例数据集。下面的示例使用的就是这个数据集；为简洁起见，相比[这里](https://ghe.clickhouse.tech/)提供的完整数据集，它采用了精简后的 schema 和部分行 (具体来说，我们只保留了与 [ClickHouse 仓库](https://github.com/ClickHouse/ClickHouse) 相关的 GitHub 事件) 。不过，这仍足以让[随该数据集发布](https://ghe.clickhouse.tech/)的大多数查询正常运行。

<div id="2-configure-clickhouse">
  ##### 2. 配置 ClickHouse
</div>

如果要连接到启用安全机制的 Kafka，则必须执行此步骤。这些设置无法通过 SQL DDL 命令传入，必须在 ClickHouse 的 config.xml 中配置。这里假设你连接的是启用了 SASL 的 instance。与 Confluent Cloud 交互时，这是最简单的方法。

```xml theme={null}
<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/` 目录下的新文件中，或合并到现有配置文件中。有关可配置的设置，请参见[此处](/zh/reference/engines/table-engines/integrations/kafka#configuration)。

我们还将创建一个名为 `KafkaEngine` 的数据库，用于本教程：

```sql theme={null}
CREATE DATABASE KafkaEngine;
```

创建数据库后，你需要切换到该数据库：

```sql theme={null}
USE KafkaEngine;
```

<div id="3-create-the-destination-table">
  ##### 3. 创建目标表
</div>

准备好目标表。为简洁起见，下面的示例使用了精简版的 GitHub schema。请注意，虽然这里使用的是 MergeTree 表引擎，但此示例也很容易改写为适用于 [MergeTree 家族](/zh/reference/engines/table-engines/mergetree-family/index) 中的任何成员。

```sql theme={null}
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)
```

<div id="4-create-and-populate-the-topic">
  ##### 4. 创建并填充 topic
</div>

接下来，我们将创建一个 topic。可以使用多种工具来完成这一步。如果 Kafka 运行在本地机器上，或在 Docker 容器中运行，[RPK](https://docs.redpanda.com/current/get-started/rpk-install/) 是一个不错的选择。我们可以运行以下命令来创建一个名为 `github`、包含 5 个分区的 topic：

```bash theme={null}
rpk topic create -p 5 github --brokers <host>:<port>
```

如果我们使用的是 Confluent Cloud 上的 Kafka，可能会更倾向于使用 [Confluent CLI](https://docs.confluent.io/platform/current/tutorials/examples/clients/docs/kcat.html#produce-records)：

```bash theme={null}
confluent kafka topic create --if-not-exists github
```

现在我们需要向这个 topic 写入一些数据，这里将使用 [kcat](https://github.com/edenhill/kcat)。如果是在本地运行 Kafka 且未启用身份验证，可以运行类似下面的命令：

```bash theme={null}
cat github_all_columns.ndjson |
kcat -P \
  -b <host>:<port> \
  -t github
```

或者，如果 Kafka 集群使用 SASL 进行身份验证，请使用以下内容：

```bash theme={null}
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 行，因此几秒内即可完成摄取。如果你想处理更大的数据集，请参阅 [ClickHouse/kafka-samples](https://github.com/ClickHouse/kafka-samples) GitHub 代码仓库中的[大数据集部分](https://github.com/ClickHouse/kafka-samples/tree/main/producer#large-datasets)。

<div id="5-create-the-kafka-table-engine">
  ##### 5. 创建 Kafka 表引擎
</div>

下面的示例创建了一个表引擎，其 schema 与 merge tree 表相同。这并非绝对必要，因为你可以在目标表中使用别名或临时列。不过，这些设置很重要——请注意，这里使用 `JSONEachRow` 作为从 Kafka topic 消费 JSON 时的数据类型。`github` 和 `clickhouse` 这两个值分别表示 topic 名称和消费者组名称。实际上，topics 也可以是一个值列表。

```sql theme={null}
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 查询后，应该能读取到一些行。请注意，这会将消费者偏移量向前移动，因此如果不进行[重置](#common-operations)，就无法再次读取这些行。另请注意其中的 limit，以及必需参数 `stream_like_engine_allow_direct_select.`

<div id="6-create-the-materialized-view">
  ##### 6. 创建 materialized view
</div>

materialized view 会连接前面创建的两个表：从 Kafka table engine 读取数据，并将其插入目标 merge tree 表。我们可以在此过程中进行多种数据转换。这里我们只做一次简单的读取和插入。使用 \* 的前提是列名完全一致 (区分大小写) 。

```sql theme={null}
CREATE MATERIALIZED VIEW github_mv TO github AS
SELECT *
FROM github_queue;
```

创建后，materialized view 会连接到 Kafka 引擎并开始读取，将行插入目标表。该过程会持续运行，后续插入到 Kafka 的消息也会被持续消费。你可以根据需要重新运行插入脚本，向 Kafka 再插入更多消息。

<div id="7-confirm-rows-have-been-inserted">
  ##### 7. 确认数据行已插入
</div>

确认目标表中已有数据：

```sql theme={null}
SELECT count() FROM github;
```

你应该会看到 200,000 行数据：

```response theme={null}
┌─count()─┐
│  200000 │
└─────────┘
```

<div id="common-operations">
  #### 常用操作
</div>

<div id="stopping--restarting-message-consumption">
  ##### 停止和恢复消息消费
</div>

要停止消息消费，可以分离 Kafka 引擎表：

```sql theme={null}
DETACH TABLE github_queue;
```

这不会影响消费者组的偏移量。要重启消费并从之前的偏移量继续，请重新附加该表。

```sql theme={null}
ATTACH TABLE github_queue;
```

<div id="adding-kafka-metadata">
  ##### 添加 Kafka 元数据
</div>

将原始 Kafka 消息中的元数据在摄取到 ClickHouse 后保留下来，通常会很有帮助。例如，我们可能希望了解某个特定 topic 或分区已经消费了多少。为此，Kafka 表引擎提供了多个[虚拟列](/zh/reference/engines/table-engines/index#table_engines-virtual_columns)。通过修改 schema 和 materialized view 的 select 语句，可以将这些虚拟列作为普通列持久化到目标表中。

首先，在向目标表添加列之前，先执行上文所述的停止操作。

```sql theme={null}
DETACH TABLE github_queue;
```

下面我们添加信息列，用于标识来源 topic 以及该行来自哪个分区。

```sql theme={null}
ALTER TABLE github
   ADD COLUMN topic String,
   ADD COLUMN partition UInt64;
```

接下来，我们需要确保虚拟列已按要求映射。
虚拟列带有 `_` 前缀。
虚拟列的完整列表可在[此处](/zh/reference/engines/table-engines/integrations/kafka#virtual-columns)查看。

要使用这些虚拟列更新表，我们需要删除 materialized view，重新 Attach Kafka 引擎表，并重新创建 materialized view。

```sql theme={null}
DROP VIEW github_mv;
```

```sql theme={null}
ATTACH TABLE github_queue;
```

```sql theme={null}
CREATE MATERIALIZED VIEW github_mv TO github AS
SELECT *, _topic AS topic, _partition as partition
FROM github_queue;
```

新读取的行应包含这些元数据。

```sql theme={null}
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         |

<div id="modify-kafka-engine-settings">
  ##### 修改 Kafka 引擎设置
</div>

我们建议删除 Kafka 引擎表，并使用新设置重新创建。在此过程中，无需修改 materialized view——Kafka 引擎表重建后，消息消费会自动恢复。

<div id="debugging-issues">
  ##### 调试问题
</div>

身份验证等错误不会出现在 Kafka 引擎 DDL 的响应中。要诊断此类问题，建议查看 ClickHouse 的主日志文件 clickhouse-server.err.log。还可以通过配置为底层 Kafka 客户端库 [librdkafka](https://github.com/edenhill/librdkafka) 启用更详细的 trace 日志。

```xml theme={null}
<kafka>
   <debug>all</debug>
</kafka>
```

<div id="handling-malformed-messages">
  ##### 处理格式错误的消息
</div>

Kafka 常常被当作数据“堆放场”使用。这会导致 topic 中混杂着不同的消息格式和不一致的字段名。应尽量避免这种情况，并利用 Kafka 的功能 (例如 Kafka Streams 或 ksqlDB) ，确保消息在写入 Kafka 之前就是格式良好且一致的。如果无法采用这些方案，ClickHouse 也提供了一些可用于缓解问题的功能。

* 将消息字段按字符串处理。如有需要，可以在 materialized view 语句中使用函数进行清洗和类型转换。虽然这不应视为生产环境方案，但对于一次性摄取可能会有帮助。
* 如果你从某个 topic 中消费 JSON，并使用 JSONEachRow format，请使用设置 [`input_format_skip_unknown_fields`](/zh/reference/settings/formats#input_format_skip_unknown_fields)。写入数据时，默认情况下，如果输入数据包含目标表中不存在的列，ClickHouse 会抛出异常。但如果启用此选项，这些多出的列会被忽略。同样，这也不是生产级方案，而且可能会让其他人感到困惑。
* 可以考虑使用设置 `kafka_skip_broken_messages`。该设置要求用户为每个块中格式错误的消息指定容忍度，并结合 `kafka_max_block_size` 来判断。如果超过这个容忍度 (按消息绝对数量计算) ，则会恢复默认的异常行为，并跳过其他消息。

<div id="delivery-semantics-and-challenges-with-duplicates">
  ##### 投递语义以及重复数据带来的挑战
</div>

Kafka 表引擎具有至少一次 (at-least-once) 投递语义。在一些已知但罕见的情况下，可能会出现重复数据。例如，消息可能已经从 Kafka 读取并成功插入 ClickHouse。但在提交新的 偏移量 之前，与 Kafka 的连接丢失了。在这种情况下，就需要重试该块。如果将分布式表或 ReplicatedMergeTree 用作目标表，则该块可以[去重](/zh/reference/engines/table-engines/mergetree-family/replication)。虽然这会降低重复行出现的概率，但它依赖于块完全一致。像 Kafka 再均衡这样的事件可能会破坏这一前提，从而在少数情况下导致重复数据。

<div id="quorum-based-inserts">
  ##### 基于仲裁的插入
</div>

在 ClickHouse 中，如果需要更高的投递保障，可能需要使用[基于仲裁的插入](/zh/reference/settings/session-settings#insert_quorum)。这项设置不能在 materialized view 或目标表上配置，但可以为用户 profile 设置，例如：

```xml theme={null}
<profiles>
  <default>
    <insert_quorum>2</insert_quorum>
  </default>
</profiles>
```

<div id="clickhouse-to-kafka">
  ### ClickHouse 到 Kafka
</div>

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

<div id="steps">
  #### 步骤
</div>

我们的初始目标如下图所示：

<Image img="https://mintcdn.com/private-7c7dfe99-detect-table-modification/j2pAbv7ihJZXp9qi/images/integrations/data-ingestion/kafka/kafka_02.webp?fit=max&auto=format&n=j2pAbv7ihJZXp9qi&q=85&s=34b92d181468a4ce0e30561c9c90a262" size="lg" alt="带有插入操作的 Kafka table engine 示意图" width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_02.webp" />

我们假设你已按照 [Kafka to ClickHouse](#kafka-to-clickhouse) 中的步骤创建好这些表和视图，并且该 topic 中的数据已被完全消费。

<div id="1-inserting-rows-directly">
  ##### 1. 直接插入行
</div>

首先，确认目标表中的行数。

```sql theme={null}
SELECT count() FROM github;
```

此时应有 200,000 行：

```response theme={null}
┌─count()─┐
│  200000 │
└─────────┘
```

现在，将 GitHub 目标表中的行重新插入到 Kafka 表引擎 github\_queue 中。请注意，这里使用的是 JSONEachRow 格式，并将 SELECT 的 LIMIT 设为 100。

```sql theme={null}
INSERT INTO github_queue SELECT * FROM github LIMIT 100 FORMAT JSONEachRow
```

重新统计 GitHub 表中的行数，确认其已增加 100。如上图所示，数据行先通过 Kafka 表引擎写入 Kafka，随后再由同一引擎重新读取，并由我们的 materialized view 插入到 GitHub 目标表中！

```sql theme={null}
SELECT count() FROM github;
```

你应该会看到新增了 100 行：

```response theme={null}
┌─count()─┐
│  200100 │
└─────────┘
```

<div id="2-using-materialized-views">
  ##### 2. 使用 materialized view
</div>

当文档插入表中时，我们可以利用 materialized view 将消息推送到 Kafka 引擎 (以及某个 topic) 。当行插入 GitHub 表时，会触发一个 materialized view，进而将这些行重新插入到 Kafka 引擎中，并写入一个新的 topic。如下图所示：

<Image img="https://mintcdn.com/private-7c7dfe99-detect-table-modification/j2pAbv7ihJZXp9qi/images/integrations/data-ingestion/kafka/kafka_03.webp?fit=max&auto=format&n=j2pAbv7ihJZXp9qi&q=85&s=996c59f123ed2bf906e0077c617af004" size="lg" alt="带有 materialized view 的 Kafka 表引擎示意图" width="2048" height="870" data-path="images/integrations/data-ingestion/kafka/kafka_03.webp" />

创建一个新的 Kafka topic `github_out` 或等效项。确保 Kafka 表引擎 `github_out_queue` 指向该 topic。

```sql theme={null}
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。

```sql theme={null}
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 to ClickHouse](#kafka-to-clickhouse) 中创建) 插入数据，文档就会自动出现在 "github\_clickhouse" topic 中。你可以使用原生 Kafka 工具来确认这一点。例如，下面我们使用 [kcat](https://github.com/edenhill/kcat) 向由 Confluent Cloud 托管的 github topic 插入 100 行数据：

```sql theme={null}
head -n 10 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>
```

读取 `github_out` topic 应可确认消息已送达。

```sql theme={null}
kcat -C \
  -b <host>:<port> \
  -t github_out \
  -X security.protocol=sasl_ssl \
  -X sasl.mechanisms=PLAIN \
  -X sasl.username=<username> \
  -X sasl.password=<password> \
  -e -q |
wc -l
```

虽然这是一个较复杂的示例，但它展示了 materialized view 与 Kafka 引擎结合使用时的强大能力。

<div id="clusters-and-performance">
  ### 集群与性能
</div>

<div id="working-with-clickhouse-clusters">
  #### 使用 ClickHouse 集群
</div>

通过 Kafka 消费者组，多个 ClickHouse 实例可以同时从同一个 topic 读取数据。每个消费者都会以 1:1 的映射关系分配到一个 topic 分区。在使用 Kafka 表引擎对 ClickHouse 的消费能力进行扩缩容时，请注意，集群中的消费者总数不能超过该 topic 的分区数。因此，请务必提前为 topic 配置好合适的分区方案。

多个 ClickHouse 实例也可以配置为使用同一个消费者组 id 从某个 topic 读取数据——该 id 在创建 Kafka 表引擎时指定。因此，每个实例都会从一个或多个分区读取数据，并将数据分段插入其本地目标表。目标表则可以进一步配置为使用 ReplicatedMergeTree 来处理数据重复。这种方法可以让 Kafka 读取能力随着 ClickHouse 集群一同扩展，前提是 Kafka 有足够多的分区。

<Image img="https://mintcdn.com/private-7c7dfe99-detect-table-modification/j2pAbv7ihJZXp9qi/images/integrations/data-ingestion/kafka/kafka_04.webp?fit=max&auto=format&n=j2pAbv7ihJZXp9qi&q=85&s=76ea6ac43a7b76ffbdc86e3e8f1a0a8f" size="lg" alt="带有 ClickHouse 集群的 Kafka 表引擎示意图" width="2048" height="1094" data-path="images/integrations/data-ingestion/kafka/kafka_04.webp" />

<div id="tuning-performance">
  #### 性能调优
</div>

在尝试提升 Kafka Engine 表的吞吐性能时，请考虑以下几点：

* 性能会因消息大小、格式以及目标表类型而异。对于单个表引擎，达到 100k 行/秒通常是可实现的。默认情况下，消息会按块读取，由参数 `kafka_max_block_size` 控制。其默认值为 [max\_insert\_block\_size](/zh/reference/settings/session-settings#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](/zh/reference/settings/server-settings/settings#background_message_broker_schedule_pool_size) 用于设置执行后台任务的线程数。这些线程会用于 Kafka 流式处理。该设置会在 ClickHouse server 启动时生效，且不能在用户 session 中更改，默认值为 16。如果你在日志中看到超时，适当增大该值可能是合适的。
* 与 Kafka 通信时使用的是 `librdkafka` 库，而它本身也会创建线程。因此，大量 Kafka 表或消费者可能会导致大量上下文切换。可以将这部分负载分散到整个集群中，并尽可能只复制目标表；或者考虑使用一个表引擎从多个 topic 读取数据——支持值列表。单个表也可以被多个 materialized view 读取，每个视图分别过滤特定 topic 的数据。

任何设置变更都应经过测试。我们建议监控 Kafka 消费者滞后，以确保扩容得当。

<div id="additional-settings">
  #### 其他设置
</div>

除了上文介绍的设置外，以下内容也值得关注：

* [Kafka\_max\_wait\_ms](/zh/reference/settings/session-settings#kafka_max_wait_ms) - 重试前从 Kafka 读取消息的等待时间，以毫秒为单位。在用户 profile 级别设置，默认值为 5000。

底层 librdkafka 的[所有设置 ](https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md)也可以放在 ClickHouse 配置文件中的 *kafka* 元素内——设置名称应写成 XML 元素，并将句点替换为下划线，例如：

```xml theme={null}
<clickhouse>
   <kafka>
       <enable_ssl_certificate_verification>false</enable_ssl_certificate_verification>
   </kafka>
</clickhouse>
```

这些是专家级设置，建议参考 Kafka 文档了解更深入的说明。
