> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/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](/docs/zh/integrations/clickpipes/home)。ClickPipes 原生支持私有网络连接，可对摄取和集群资源分别进行扩缩容，并为流式 Kafka 数据摄取到 ClickHouse 提供全面监控。
</Note>

要使用 Kafka 表引擎，您应当对 [ClickHouse materialized views](/docs/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 家族](/docs/zh/reference/engines/table-engines/mergetree-family/index) 的表。如下图所示：

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_01.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=fd6990e133e6b46eb62f7057956fb5a3" 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>

<Steps>
  <Step title="准备" id="1-prepare">
    如果你的目标 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/)的大多数查询正常运行。
  </Step>

  <Step title="配置 ClickHouse" id="2-configure-clickhouse">
    如果您要连接到启用了安全机制的 Kafka，则此步骤必不可少。这些设置无法通过 SQL DDL 命令传入，必须在 ClickHouse 的 config.xml 中配置。这里假设您连接的是启用了 SASL 保护的实例。这是在与 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/` 目录下的新文件中，或者将其合并到现有配置文件中。有关可配置的设置，请参见[此处](/docs/zh/reference/engines/table-engines/integrations/kafka#configuration)。

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

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

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

    ```sql theme={null}
    USE KafkaEngine;
    ```
  </Step>

  <Step title="创建目标表" id="3-create-the-destination-table">
    准备好目标表。为简洁起见，下面的示例使用了精简版的 GitHub schema。请注意，虽然这里使用的是 MergeTree 表引擎，但此示例也很容易改写为适用于 [MergeTree 家族](/docs/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)
    ```
  </Step>

  <Step title="创建 topic 并写入数据" id="4-create-and-populate-the-topic">
    接下来，我们将创建一个 topic。为此可以使用多种工具。如果是在本机上本地运行 Kafka，或者在 Docker 容器中运行 Kafka，[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 命令行客户端](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 行，因此只需几秒钟即可完成摄取。如果你想使用更大的数据集，请查看 GitHub 仓库 [ClickHouse/kafka-samples](https://github.com/ClickHouse/kafka-samples) 中的[大数据集部分](https://github.com/ClickHouse/kafka-samples/tree/main/producer#large-datasets)。
  </Step>

  <Step title="创建 Kafka 表引擎" id="5-create-the-kafka-table-engine">
    下面的示例创建了一个与 MergeTree 表具有相同 schema 的表引擎。这并非严格必需，因为你可以在目标表中使用别名列或临时列。不过，这些设置很重要——请注意，这里使用 `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 应该能读出一些行。请注意，这会将消费者 offsets 向前推进，因此如果不进行[重置](#common-operations)，这些行将无法再次读取。另请注意限制以及必需参数 `stream_like_engine_allow_direct_select.`
  </Step>

  <Step title="创建 materialized view" id="6-create-the-materialized-view">
    materialized view 会连接前面创建的两个表，从 Kafka 表引擎读取数据，并将其插入目标 MergeTree 表。我们可以进行多种数据转换。这里我们将执行简单的读取和插入操作。使用 \* 的前提是列名完全一致 (区分大小写) 。

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

    在创建时，materialized view 会连接到 Kafka 引擎并开始读取，将行插入目标表。此过程会无限期持续，之后插入到 Kafka 的消息也会被消费。你可以随时重新运行插入脚本，向 Kafka 再插入更多消息。
  </Step>

  <Step title="确认行已插入" id="7-confirm-rows-have-been-inserted">
    确认目标表中有数据：

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

    你应该能看到 200,000 行：

    ```response theme={null}
    ┌─count()─┐
    │  200000 │
    └─────────┘
    ```
  </Step>
</Steps>

<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 表引擎提供了多个[虚拟列](/docs/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;
```

接下来，我们需要确保虚拟列已按要求映射。
虚拟列带有 `_` 前缀。
虚拟列的完整列表可在[此处](/docs/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`](/docs/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 用作目标表，则该块可以[去重](/docs/zh/reference/engines/table-engines/mergetree-family/replication)。虽然这会降低重复行出现的概率，但它依赖于块完全一致。像 Kafka 再均衡这样的事件可能会破坏这一前提，从而在少数情况下导致重复数据。

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

在 ClickHouse 中，如果需要更高的投递保障，可能需要使用[基于仲裁的插入](/docs/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-1">
  #### 步骤
</div>

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

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

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

<Steps>
  <Step title="直接插入数据行" id="1-inserting-rows-directly">
    首先，确认目标表中的行数。

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

    此时应有 200,000 行：

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

    现在将行从 GitHub 目标表重新插入到 Kafka 表引擎 github\_queue 中。请注意，我们使用了 JSONEachRow 格式，并通过 LIMIT 将 select 限制为 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 │
    └─────────┘
    ```
  </Step>

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

    <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_03.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ead2a8c956700436a1fd15bb803ef49b" 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 到 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 引擎结合使用时的强大能力。
  </Step>
</Steps>

<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/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/kafka/kafka_04.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=fc14aa743c7b5cda9b7d722e18bc59f0" 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](/docs/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](/docs/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](/docs/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 文档了解更深入的说明。
