> ## 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/ja/integrations/clickpipes/home) の利用をお勧めします。ClickPipes は、プライベートネットワーク接続に対応しており、インジェストとクラスターリソースを個別にスケーリングできるほか、Kafka データを ClickHouse にストリーミングで取り込むための包括的な監視機能を備えています。
</Note>

Kafka テーブルエンジン を使用するには、[ClickHouse materialized views](/docs/ja/concepts/features/materialized-views/cascading-materialized-views) について一通り理解している必要があります。

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

まずは、最も一般的なユースケースである、Kafka テーブルエンジン を使って Kafka から ClickHouse にデータを挿入する方法を見ていきます。

Kafka テーブルエンジン を使うと、ClickHouse は Kafka トピックから直接データを読み取れます。トピック上のメッセージを確認する用途には便利ですが、この engine は設計上、一回限りの取得しかできません。つまり、テーブルエンジン に対してクエリを実行すると、結果を呼び出し元に返す前にキューからデータを消費し、consumer の OFFSET を進めます。そのため、これらの OFFSET をリセットしない限り、実質的にデータを再読込することはできません。

テーブルエンジン から読み取ったデータを永続化するには、そのデータを取り込み、別の table に挿入する仕組みが必要です。トリガーベースの materialized view は、この機能をネイティブに提供します。materialized view は テーブルエンジン に対する読み取りを開始し、ドキュメントの batches を受け取ります。TO 句はデータの宛先を決定します。通常は [MergeTree ファミリー](/docs/ja/reference/engines/table-engines/mergetree-family/index) の table です。この処理を以下に示します。

<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">
    ターゲットトピックにデータが投入されている場合は、以下の内容を調整して自分のデータセットに使用できます。あるいは、サンプルの GitHub データセットを[こちら](https://datasets-documentation.s3.eu-west-3.amazonaws.com/kafka/github_all_columns.ndjson)で提供しています。このデータセットは以下の例で使用しており、簡潔にするため、[こちら](https://ghe.clickhouse.tech/)で利用できる完全なデータセットと比べて、簡略化したスキーマと行の一部のみ (具体的には、[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/ja/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 のスキーマを使用しています。ここでは MergeTree テーブルエンジンを使用していますが、この例は [MergeTree family](/docs/ja/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="トピックを作成し、データを投入する" id="4-create-and-populate-the-topic">
    次に、トピックを作成します。これを行うために使用できるツールはいくつかあります。Kafka をローカルマシン上または Docker コンテナ内で実行している場合は、[RPK](https://docs.redpanda.com/current/get-started/rpk-install/) が適しています。以下のコマンドを実行すると、5 つのパーティションを持つ `github` という名前のトピックを作成できます。

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

    Kafka を Confluent Cloud 上で実行している場合は、[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 をローカルで実行しており、authentication が無効になっている場合は、以下のようなコマンドを実行できます。

    ```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 リポジトリの [large datasets section](https://github.com/ClickHouse/kafka-samples/tree/main/producer#large-datasets) を参照してください。
  </Step>

  <Step title="Kafkaテーブルエンジンを作成する" id="5-create-the-kafka-table-engine">
    以下の例では、MergeTree テーブルと同じスキーマを持つテーブルエンジンを作成します。これは厳密には必須ではなく、ターゲットテーブルにエイリアスや一時的なカラムを含めることもできます。ただし、設定は重要です。Kafka トピックから JSON を取り込むデータ型として `JSONEachRow` を使用している点に注意してください。値 `github` と `clickhouse` は、それぞれトピック名とコンシューマグループ名を表します。実際には、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 でいくつかの行を読み取れるはずです。  これによりコンシューマーオフセットが進むため、[reset](#common-operations) を行わない限り、これらの行は再度読み取れなくなる点に注意してください。制限事項と必須パラメータ `stream_like_engine_allow_direct_select` に注意してください。
  </Step>

  <Step title="materialized viewを作成する" id="6-create-the-materialized-view">
    materialized view は、前に作成した 2 つのテーブルを接続し、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>

ClickHouse に取り込んだ後も、元の Kafka メッセージのメタデータを追跡できると便利です。たとえば、特定のトピックやパーティションをどれだけ消費したかを把握したい場合があります。このため、Kafka テーブルエンジンは複数の[仮想カラム](/docs/ja/reference/engines/table-engines/index#table_engines-virtual_columns)を公開しています。スキーマと materialized view の SELECT ステートメントを変更することで、これらをターゲットテーブルのカラムとして永続化できます。

まず、ターゲットテーブルにカラムを追加する前に、前述の停止操作を実行します。

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

以下では、各行の取得元のトピックとパーティションを識別するための情報カラムを追加します。

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

次に、必要な仮想カラムが適切にマッピングされていることを確認する必要があります。
仮想カラムは `_` で始まります。
仮想カラムの完全な一覧は[こちら](/docs/ja/reference/engines/table-engines/integrations/kafka#virtual-columns)で確認できます。

仮想カラムをテーブルに反映するには、materialized view を削除し、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) については、設定によりさらに詳細なトレースログを有効にできます。

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

<div id="handling-malformed-messages">
  ##### 不正な形式のメッセージの処理
</div>

Kafka は、しばしばデータの「投げ込み先」として使われます。その結果、トピック内で複数のメッセージ形式が混在したり、フィールド名に一貫性がなくなったりします。こうした状況は避け、Kafka に書き込む前にメッセージが整形式かつ一貫したものになるよう、Kafka Streams や ksqlDB などの Kafka の機能を活用してください。これらの方法が使えない場合でも、ClickHouse には対処に役立つ機能がいくつかあります。

* メッセージのフィールドは文字列として扱います。必要に応じて、materialized view ステートメント内で関数を使ってクレンジングや CAST を行えます。これは本番向けの解決策ではありませんが、一度限りのインジェストには役立つ場合があります。
* トピックから JSON を読み込み、JSONEachRow フォーマットを使用している場合は、設定 [`input_format_skip_unknown_fields`](/docs/ja/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 テーブルエンジン は少なくとも 1 回の配信セマンティクスを持っています。既知のまれな状況がいくつかあり、その場合は重複が発生する可能性があります。たとえば、メッセージが Kafka から読み取られ、ClickHouse への挿入に成功することがあります。新しいオフセットをコミットする前に、Kafka への接続が失われる可能性があります。この状況では、ブロックの再試行が必要です。ブロックは、ターゲットテーブルとして分散テーブルまたは ReplicatedMergeTree を使用することで[重複排除](/docs/ja/reference/engines/table-engines/mergetree-family/replication)できます。これにより重複する行が発生する可能性は低くなりますが、これはブロックが同一であることを前提としています。Kafka のリバランシングのような事象によってこの前提が崩れ、まれに重複が発生することがあります。

<div id="quorum-based-inserts">
  ##### クォーラムベースのインサート
</div>

ClickHouse でより高い配信保証が必要な場合は、[クォーラムベースのインサート](/docs/ja/reference/settings/session-settings#insert_quorum) が必要になることがあります。これは materialized view やターゲットテーブルには設定できませんが、ユーザープロファイルには設定できます。たとえば次のようになります。

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

<div id="clickhouse-to-kafka">
  ### ClickHouse から Kafka へ
</div>

比較的まれなユースケースではありますが、ClickHouse のデータを Kafka に永続化することもできます。たとえば、ここでは Kafka テーブルエンジンに手動で行を insert します。このデータは同じ Kafka エンジンによって読み取られ、その materialized view によってデータが MergeTree テーブルに格納されます。最後に、既存の source table からテーブルを読み取るために、Kafka への insert で materialized view を適用する方法を示します。

<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="insert を伴う Kafka テーブルエンジンの図" width="2048" height="852" data-path="images/integrations/data-ingestion/kafka/kafka_02.webp" />

[Kafka to ClickHouse](#kafka-to-clickhouse) の手順で作成したテーブルとビューが存在し、トピックは完全に消費済みであることを前提とします。

<Steps>
  <Step title="行を直接 insert する" id="1-inserting-rows-directly">
    まず、ターゲットテーブルの行数を確認してください。

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

    200,000行あるはずです。

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

    次に、GitHub ターゲットテーブルから行を Kafka テーブルエンジン github\_queue に再度 insert します。JSONEachRow フォーマットを利用し、SELECT を 100 行に LIMIT している点に注目してください。

    ```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 エンジン (およびトピック) に送ることができます。GitHub テーブルに行が挿入されると、materialized view がトリガーされ、その結果、行は Kafka エンジン に再度挿入され、新しいトピックにも送られます。これも図で見るのが最もわかりやすいでしょう。

    <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 トピック `github_out` または同等のものを作成してください。Kafka テーブルエンジン `github_out_queue` がこのトピックを指していることを確認してください。

    ```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 テーブルを参照するように設定します。これがトリガーされると、上記のエンジンに行が insert されます。これにより、GitHub テーブルへの追加分は新しい Kafka トピック にプッシュされます。

    ```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トピック ([Kafka to ClickHouse](#kafka-to-clickhouse) の一部として作成) にinsertすると、ドキュメントは自動的に "github\_clickhouse" トピックに現れます。これをKafkaのネイティブツールで確認してください。たとえば以下では、Confluent Cloudでホストされているトピックに対して、[kcat](https://github.com/edenhill/kcat) を使用し、githubトピックに100行をinsertします。

    ```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` トピックを読めば、メッセージが配信されたことを確認できるはずです。

    ```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
    ```

    これは複雑な例ですが、Kafka エンジンと組み合わせて使用した場合の materialized view の威力を示しています。
  </Step>
</Steps>

<div id="clusters-and-performance">
  ### クラスターとパフォーマンス
</div>

<div id="working-with-clickhouse-clusters">
  #### ClickHouseクラスターの使用
</div>

Kafkaのコンシューマグループを使うことで、複数のClickHouseインスタンスが同じトピックを読み取ることができます。各コンシューマーは、トピック内の1つのパーティションに1:1で割り当てられます。Kafka テーブルエンジンを使用してClickHouseのデータ取り込みをスケールする場合、クラスター内のコンシューマー総数はトピックのパーティション数を超えられない点に注意してください。そのため、対象のトピックでは事前に適切なパーティション化を設定しておく必要があります。

複数のClickHouseインスタンスを、同じ コンシューマグループ id を使って1つのトピックを読み取るように設定できます。これはKafka テーブルエンジンの作成時に指定します。その結果、各インスタンスは1つ以上のパーティションから読み取り、ローカルのターゲットテーブルにセグメントを挿入します。さらに、ターゲットテーブルは、データの重複を処理するために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 テーブルのスループットを向上させる際は、次の点を考慮してください。

* パフォーマンスは、メッセージサイズ、フォーマット、ターゲットテーブルの types によって異なります。単一のテーブルエンジンで 10 万行/秒は十分に達成可能な水準です。デフォルトでは、メッセージは `kafka_max_block_size` パラメーターで制御されるブロック単位で読み取られます。これはデフォルトで [max\_insert\_block\_size](/docs/ja/reference/settings/session-settings#max_insert_block_size) に設定されており、既定値は 1,048,576 です。メッセージが極端に大きい場合を除き、この値はほぼ常に増やすべきです。500k ～ 1M 程度の値も珍しくありません。スループットへの影響をテストして評価してください。
* テーブルエンジンのコンシューマー数は `kafka_num_consumers` で増やせます。ただし、デフォルトでは `kafka_thread_per_consumer` を既定値の `1` から変更しない限り、インサートは単一スレッドで直列化されます。フラッシュが並列に実行されるよう、これを `1` に設定してください。なお、N 個のコンシューマーを持つ Kafka engine テーブル (`kafka_thread_per_consumer=1`) を作成することは、それぞれに materialized view があり、`kafka_thread_per_consumer=0` が設定された N 個の Kafka engine を作成するのと論理的に等価です。
* コンシューマー数の増加は無償ではありません。各コンシューマーは独自のバッファーとスレッドを保持するため、サーバーのオーバーヘッドが増えます。まずは可能であればクラスター全体で線形にスケールさせ、コンシューマーによるオーバーヘッドを意識してください。
* Kafka メッセージのスループットにばらつきがあり、遅延を許容できる場合は、より大きなブロックがフラッシュされるよう `stream_flush_interval_ms` を増やすことを検討してください。
* [background\_message\_broker\_schedule\_pool\_size](/docs/ja/reference/settings/server-settings/settings#background_message_broker_schedule_pool_size) は、バックグラウンドタスクを実行するスレッド数を設定します。これらのスレッドは Kafka ストリーミングに使用されます。この設定は ClickHouseサーバーの起動時に適用され、ユーザーセッション中に変更することはできません。既定値は 16 です。ログにタイムアウトが見られる場合は、この値を増やすのが適切なことがあります。
* Kafka との通信には `librdkafka` ライブラリが使用されており、このライブラリ自体もスレッドを作成します。そのため、多数の Kafka テーブルやコンシューマーがあると、大量のコンテキストスイッチが発生する可能性があります。この負荷はクラスター全体に分散し、可能であればターゲットテーブルのみをレプリケートするか、1 つのテーブルエンジンで複数のトピックを読み取ることを検討してください。値のリストがサポートされています。1 つのテーブルから複数の materialized view を読み取ることができ、それぞれが特定のトピックのデータをフィルタリングできます。

設定の変更は必ずテストしてください。適切にスケールできていることを確認するため、Kafka のコンシューマラグを監視することを推奨します。

<div id="additional-settings">
  #### 追加の設定
</div>

上述の設定に加えて、以下の設定も参考になります。

* [Kafka\_max\_wait\_ms](/docs/ja/reference/settings/session-settings#kafka_max_wait_ms) - 再試行する前に Kafka からメッセージを読み取る際の待機時間 (ミリ秒) です。ユーザープロファイルレベルで設定し、デフォルト値は 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 のドキュメントを参照することをお勧めします。
