メインコンテンツへスキップ
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 トピックから直接データを読み取れます。トピック上のメッセージを確認する用途には便利ですが、この engine は設計上、一回限りの取得しかできません。つまり、table に対してクエリを実行すると、結果を呼び出し元に返す前にキューからデータを消費し、consumer の OFFSET を進めます。そのため、これらの OFFSET をリセットしない限り、実質的にデータを再読込することはできません。 table engine から読み取ったデータを永続化するには、そのデータを取り込み、別の table に挿入する仕組みが必要です。トリガーベースの materialized view は、この機能をネイティブに提供します。materialized view は table engine に対する読み取りを開始し、ドキュメントの batches を受け取ります。TO 句はデータの宛先を決定します。通常は MergeTree ファミリー の table です。この処理を以下に示します。

手順

1. 準備
対象のトピックにデータが投入されている場合は、以下の内容を自身のデータセットに合わせて調整して使用できます。あるいは、サンプルの GitHub データセットをこちらから利用できます。このデータセットは以下の例で使用しており、簡潔にするため、完全なデータセット (こちらで利用可能) に比べて、スキーマを簡略化し、行の一部のみを含んでいます (具体的には、ClickHouse リポジトリに関する GitHub イベントに限定しています) 。それでも、このデータセットとともに公開されているクエリの大半を実行するには十分です。
2. ClickHouse を設定する
この手順は、セキュアな Kafka に接続する場合に必要です。これらの設定は SQL DDL コマンドでは指定できないため、ClickHouse の config.xml で設定する必要があります。ここでは、SASL で保護されたインスタンスに接続することを前提としています。これは、Confluent Cloud と連携する際に最も簡単な方法です。
上記のスニペットは、conf.d/ ディレクトリ配下の新しいファイルに配置するか、既存の設定ファイルに統合してください。設定可能な項目については、こちらを参照してください。 このチュートリアルで使用するため、KafkaEngine という名前のデータベースも作成します:
データベースを作成したら、そのデータベースに切り替えます:
3. 宛先テーブルを作成する
宛先テーブルを準備します。以下の例では、簡潔にするため、簡略化した GitHub のスキーマを使用しています。ここでは MergeTree テーブルエンジンを使用していますが、この例は MergeTree family のどのメンバーにも簡単に適用できます。
4. トピックを作成してデータを投入する
次に、トピックを作成します。これには、いくつかのツールを使用できます。Kafka をローカルマシン上または Docker コンテナー内で実行している場合は、RPK が便利です。次のコマンドを実行すると、github という名前のトピックを 5 つのパーティションで作成できます。
Kafka を Confluent Cloud 上で実行している場合は、Confluent CLI を使用するのがよいでしょう。
次に、このトピックにデータをいくつか投入する必要があります。これには kcat を使用します。Kafka をローカルで実行し、認証を無効にしている場合は、次のようなコマンドを実行できます。
または、Kafkaクラスターで認証にSASLを使用している場合は、以下を使用します。
このデータセットには 200,000 行が含まれているため、取り込みは数秒で完了するはずです。より大きなデータセットを扱いたい場合は、ClickHouse/kafka-samples GitHub リポジトリのlarge datasets セクションを参照してください。
5. Kafka テーブルエンジン を作成する
以下の例では、MergeTree テーブルと同じスキーマを持つ table engine を作成します。これは厳密には必須ではなく、ターゲットテーブルに alias や一時的なカラムを含めることもできます。ただし、設定は重要です。特に、Kafka トピック から JSON を取り込むためのデータ型として JSONEachRow を使用している点に注意してください。githubclickhouse という値は、それぞれ トピック 名とコンシューマグループ名を表します。トピックs には実際には複数の値を指定できます。
以下では、エンジンの設定とパフォーマンスチューニングについて説明します。この時点で、テーブル github_queue に対して単純な select を実行すると、いくつかの行を読み取れるはずです。これにより consumer の offset が進むため、リセットしない限り、これらの行は再度読み取れなくなる点に注意してください。limit と、必須パラメータである stream_like_engine_allow_direct_select. にも注意してください。
6. materialized view を作成する
materialized view は、前の手順で作成した 2 つのテーブルを接続し、Kafka table engine からデータを読み取って、対象の merge tree テーブルに挿入します。データ変換はいくつか行えますが、ここでは単純な読み取りと挿入のみを行います。* の使用は、カラム名が同一であることを前提としています (大文字と小文字は区別されます) 。
作成されると、materialized view は Kafka エンジン に接続し、読み取りを開始してターゲットテーブルに行を挿入します。この処理は継続的に行われ、以降 Kafka に挿入されるメッセージも消費されます。必要に応じて挿入スクリプトを再実行し、Kafka にさらにメッセージを挿入してください。
7. 行が挿入されていることを確認する
ターゲットテーブルにデータが存在することを確認します。
200,000行と表示されるはずです。

主な操作

メッセージ消費の停止と再開
メッセージ消費を停止するには、Kafkaエンジンのテーブルをデタッチします:
これはコンシューマグループのオフセットには影響しません。コンシュームを再開して前回のオフセットから処理を続けるには、テーブルを再アタッチします。
Kafka メタデータの追加
ClickHouse に取り込んだ後も、元の Kafka メッセージのメタデータを追跡できると便利です。たとえば、特定のトピックやパーティションをどれだけ消費したかを把握したい場合があります。このため、Kafka テーブルエンジンは複数の仮想カラムを公開しています。スキーマと materialized view の SELECT ステートメントを変更することで、これらをターゲットテーブルのカラムとして永続化できます。 まず、ターゲットテーブルにカラムを追加する前に、前述の停止操作を実行します。
以下では、各行の取得元のトピックとパーティションを識別するための情報カラムを追加します。
次に、必要な仮想カラムが適切にマッピングされていることを確認する必要があります。 仮想カラムは _ で始まります。 仮想カラムの完全な一覧はこちらで確認できます。 仮想カラムをテーブルに反映するには、materialized view を削除し、Kafka エンジンのテーブルを再アタッチしてから、materialized view を再作成する必要があります。
新たに取り込まれた行には、そのメタデータが含まれているはずです。
結果は以下のようになります。
Kafka エンジン設定の変更
Kafka エンジンテーブルはいったん削除し、新しい設定で再作成することを推奨します。この作業中にmaterialized viewを変更する必要はありません。Kafka エンジンテーブルを再作成すると、メッセージの消費が再開されます。
問題のデバッグ
認証の問題などのエラーは、KafkaエンジンのDDLに対する応答には返されません。問題を診断するには、メインのClickHouseログファイル clickhouse-server.err.log を使用することを推奨します。基盤となるKafkaクライアントライブラリ librdkafka については、設定によりさらに詳細なトレースログを有効にできます。
不正な形式のメッセージの処理
Kafka は、しばしばデータの「投げ込み先」として使われます。その結果、トピック内で複数のメッセージ形式が混在したり、フィールド名に一貫性がなくなったりします。こうした状況は避け、Kafka に書き込む前にメッセージが整形式かつ一貫したものになるよう、Kafka Streams や ksqlDB などの Kafka の機能を活用してください。これらの方法が使えない場合でも、ClickHouse には対処に役立つ機能がいくつかあります。
  • メッセージのフィールドは文字列として扱います。必要に応じて、materialized view ステートメント内で関数を使ってクレンジングや CAST を行えます。これは本番向けの解決策ではありませんが、一度限りのインジェストには役立つ場合があります。
  • トピックから JSON を読み込み、JSONEachRow フォーマットを使用している場合は、設定 input_format_skip_unknown_fields を使用してください。データの書き込み時、デフォルトでは、入力データにターゲットテーブルに存在しないカラムが含まれていると、ClickHouse は例外をスローします。ただし、このオプションを有効にすると、こうした余分なカラムは無視されます。これも本番レベルの解決策ではなく、他の利用者を混乱させる可能性があります。
  • 設定 kafka_skip_broken_messages の利用も検討してください。これを使うには、不正な形式のメッセージに対するブロックごとの許容度を、kafka_max_block_size を踏まえてユーザーが指定する必要があります。この許容度を超えた場合 (絶対メッセージ数で判定) 、通常どおり例外が発生し、それ以外のメッセージはスキップされます。
配信セマンティクスと重複に関する課題
Kafka テーブルエンジン は少なくとも 1 回の配信セマンティクスを持っています。既知のまれな状況がいくつかあり、その場合は重複が発生する可能性があります。たとえば、メッセージが Kafka から読み取られ、ClickHouse への挿入に成功することがあります。新しいオフセットをコミットする前に、Kafka への接続が失われる可能性があります。この状況では、ブロックの再試行が必要です。ブロックは、ターゲットテーブルとして分散テーブルまたは ReplicatedMergeTree を使用することで重複排除できます。これにより重複する行が発生する可能性は低くなりますが、これはブロックが同一であることを前提としています。Kafka のリバランシングのような事象によってこの前提が崩れ、まれに重複が発生することがあります。
クォーラムベースのインサート
ClickHouse でより高い配信保証が必要な場合は、クォーラムベースのインサート が必要になることがあります。これは materialized view やターゲットテーブルには設定できませんが、ユーザープロファイルには設定できます。たとえば次のようになります。

ClickHouse から Kafka へ

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

手順

最初の目的は、次の図を見るとわかりやすいでしょう。 Kafka to ClickHouse の手順で作成したテーブルとビューが存在し、トピックは完全に消費済みであることを前提とします。
1. 行を直接挿入する
まず、ターゲットテーブルの行数を確認します。
200,000行になっているはずです。
ここで、GitHub のターゲットテーブルから Kafka テーブルエンジン github_queue に行を再度挿入します。JSONEachRow フォーマットを使用し、SELECT を 100 に LIMIT している点に注目してください。
GitHub 内の行数を再度数え、100 行増えていることを確認します。上の図に示すように、行は Kafka table engine を介して Kafka に挿入された後、同じ Kafka table engine で再び読み取られ、materialized view によって GitHub のターゲットテーブルに挿入されます!
さらに100行表示されるはずです。
2. materialized view を使う
テーブルにドキュメントが挿入されたとき、materialized view を利用してメッセージを Kafka エンジン (およびトピック) に送ることができます。GitHub テーブルに行が挿入されると、materialized view がトリガーされ、その結果、行は Kafka エンジン に再度挿入され、新しいトピックにも送られます。これも図で見るのが最もわかりやすいでしょう。 新しい Kafka トピック github_out または同等のものを作成してください。Kafka テーブルエンジン github_out_queue がこのトピックを指していることを確認してください。
次に、新しい materialized view github_out_mv を作成し、GitHub テーブルを参照するように設定します。これがトリガーされると、上記のエンジンに行が insert されます。これにより、GitHub テーブルへの追加分は新しい Kafka トピック にプッシュされます。
元のgithubトピック (Kafka to ClickHouse の一部として作成) にinsertすると、ドキュメントは自動的に “github_clickhouse” トピックに現れます。これをKafkaのネイティブツールで確認してください。たとえば以下では、Confluent Cloudでホストされているトピックに対して、kcat を使用し、githubトピックに100行をinsertします。
github_out トピックを読み取れば、メッセージが配信されたことを確認できるはずです。
やや複雑な例ではありますが、これはmaterialized viewをKafka エンジンと組み合わせて使用した際の威力を示しています。

クラスターとパフォーマンス

ClickHouseクラスターの使用

Kafkaのコンシューマグループを使うことで、複数のClickHouseインスタンスが同じトピックを読み取ることができます。各コンシューマーは、トピック内の1つのパーティションに1:1で割り当てられます。Kafka テーブルエンジンを使用してClickHouseのデータ取り込みをスケールする場合、クラスター内のコンシューマー総数はトピックのパーティション数を超えられない点に注意してください。そのため、対象のトピックでは事前に適切なパーティション化を設定しておく必要があります。 複数のClickHouseインスタンスを、同じ コンシューマグループ id を使って1つのトピックを読み取るように設定できます。これはKafka テーブルエンジンの作成時に指定します。その結果、各インスタンスは1つ以上のパーティションから読み取り、ローカルのターゲットテーブルにセグメントを挿入します。さらに、ターゲットテーブルは、データの重複を処理するためにReplicatedMergeTreeを使用するよう設定できます。この方法では、Kafkaのパーティション数が十分にあれば、ClickHouseクラスターに合わせてKafkaからの読み取りをスケールできます。

パフォーマンスのチューニング

Kafka Engine テーブルのスループットを向上させる際は、次の点を考慮してください。
  • パフォーマンスは、メッセージサイズ、フォーマット、ターゲットテーブルの types によって異なります。単一のテーブルエンジンで 10 万行/秒は十分に達成可能な水準です。デフォルトでは、メッセージは kafka_max_block_size パラメーターで制御されるブロック単位で読み取られます。これはデフォルトで 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 は、バックグラウンドタスクを実行するスレッド数を設定します。これらのスレッドは Kafka ストリーミングに使用されます。この設定は ClickHouseサーバーの起動時に適用され、ユーザーセッション中に変更することはできません。既定値は 16 です。ログにタイムアウトが見られる場合は、この値を増やすのが適切なことがあります。
  • Kafka との通信には librdkafka ライブラリが使用されており、このライブラリ自体もスレッドを作成します。そのため、多数の Kafka テーブルやコンシューマーがあると、大量のコンテキストスイッチが発生する可能性があります。この負荷はクラスター全体に分散し、可能であればターゲットテーブルのみをレプリケートするか、1 つのテーブルエンジンで複数のトピックを読み取ることを検討してください。値のリストがサポートされています。1 つのテーブルから複数の materialized view を読み取ることができ、それぞれが特定のトピックのデータをフィルタリングできます。
設定の変更は必ずテストしてください。適切にスケールできていることを確認するため、Kafka のコンシューマラグを監視することを推奨します。

追加の設定

上述の設定に加えて、以下の設定も参考になります。
  • Kafka_max_wait_ms - 再試行する前に Kafka からメッセージを読み取る際の待機時間 (ミリ秒) です。ユーザープロファイルレベルで設定し、デフォルト値は 5000 です。
また、基盤となる librdkafka のすべての設定 は、ClickHouse の設定ファイル内の kafka 要素にも指定できます。設定名は、ピリオドをアンダースコアに置き換えた XML 要素にする必要があります (例:) 。
これらは高度な設定のため、詳しくは Kafka のドキュメントを参照することをお勧めします。
最終更新日 2026年6月19日