はじめに
このガイドでは、named collections を使用して ClickHouse を Kafka に接続する方法を紹介します。named collections の設定ファイルを使用することで、次のような利点があります。
- 設定を一元化し、より簡単に管理できます。
- SQL のテーブル定義を変更することなく、設定を変更できます。
- 単一の設定ファイルを確認するだけで、設定のレビューやトラブルシューティングをより容易に行えます。
このガイドは、Apache Kafka 3.4.1 および ClickHouse 24.5.1 で検証されています。
前提
このドキュメントでは、以下を前提としています。
- 稼働中のKafkaクラスターがあること。
- ClickHouseクラスターが構成済みで、稼働していること。
- SQLの基本的な知識があり、ClickHouseおよびKafkaの設定に精通していること。
前提条件
named collection を作成するユーザーに必要なアクセス権限があることを確認してください。
<access_management>1</access_management>
<named_collection_control>1</named_collection_control>
<show_named_collections>1</show_named_collections>
<show_named_collections_secrets>1</show_named_collections_secrets>アクセス制御を有効にする方法の詳細については、ユーザー管理ガイドを参照してください。
設定
以下のセクションを ClickHouse の config.xml ファイルに追加します。
<!-- Kafkaインテグレーション用の名前付きコレクション -->
<named_collections>
<cluster_1>
<!-- ClickHouse Kafkaエンジンパラメータ -->
<kafka_broker_list>c1-kafka-1:9094,c1-kafka-2:9094,c1-kafka-3:9094</kafka_broker_list>
<kafka_topic_list>cluster_1_clickhouse_topic</kafka_topic_list>
<kafka_group_name>cluster_1_clickhouse_consumer</kafka_group_name>
<kafka_format>JSONEachRow</kafka_format>
<kafka_commit_every_batch>0</kafka_commit_every_batch>
<kafka_num_consumers>1</kafka_num_consumers>
<kafka_thread_per_consumer>1</kafka_thread_per_consumer>
<!-- Kafka拡張設定 -->
<kafka>
<security_protocol>SASL_SSL</security_protocol>
<enable_ssl_certificate_verification>false</enable_ssl_certificate_verification>
<sasl_mechanism>PLAIN</sasl_mechanism>
<sasl_username>kafka-client</sasl_username>
<sasl_password>kafkapassword1</sasl_password>
<debug>all</debug>
<auto_offset_reset>latest</auto_offset_reset>
</kafka>
</cluster_1>
<cluster_2>
<!-- ClickHouse Kafkaエンジンパラメータ -->
<kafka_broker_list>c2-kafka-1:29094,c2-kafka-2:29094,c2-kafka-3:29094</kafka_broker_list>
<kafka_topic_list>cluster_2_clickhouse_topic</kafka_topic_list>
<kafka_group_name>cluster_2_clickhouse_consumer</kafka_group_name>
<kafka_format>JSONEachRow</kafka_format>
<kafka_commit_every_batch>0</kafka_commit_every_batch>
<kafka_num_consumers>1</kafka_num_consumers>
<kafka_thread_per_consumer>1</kafka_thread_per_consumer>
<!-- Kafka拡張設定 -->
<kafka>
<security_protocol>SASL_SSL</security_protocol>
<enable_ssl_certificate_verification>false</enable_ssl_certificate_verification>
<sasl_mechanism>PLAIN</sasl_mechanism>
<sasl_username>kafka-client</sasl_username>
<sasl_password>kafkapassword2</sasl_password>
<debug>all</debug>
<auto_offset_reset>latest</auto_offset_reset>
</kafka>
</cluster_2>
</named_collections>設定に関する注意事項
- Kafka のアドレスと関連設定は、ご使用の Kafka クラスター構成に合わせて調整してください。
<kafka>の前のセクションには、ClickHouse の Kafka エンジンパラメータが含まれています。パラメータの一覧については、Kafka エンジンパラメータ を参照してください。<kafka>内のセクションには、Kafka の拡張設定オプションが含まれています。その他のオプションについては、librdkafka の設定を参照してください。- この例では、
SASL_SSLセキュリティプロトコルとPLAINメカニズムを使用しています。これらの設定は、Kafka クラスターの構成に応じて調整してください。
テーブルとデータベースの作成
ClickHouseクラスター上に必要なデータベースとテーブルを作成します。ClickHouseを単一ノードで実行している場合は、SQLコマンドのクラスター指定部分を省略し、ReplicatedMergeTree の代わりに別のエンジンを使用してください。
データベースを作成する
CREATE DATABASE kafka_testing ON CLUSTER LAB_CLICKHOUSE_CLUSTER;Kafkaテーブルを作成する
1つ目のKafkaクラスター用に、1つ目のKafkaテーブルを作成します。
CREATE TABLE kafka_testing.first_kafka_table ON CLUSTER LAB_CLICKHOUSE_CLUSTER
(
`id` UInt32,
`first_name` String,
`last_name` String
)
ENGINE = Kafka(cluster_1);2 つ目の Kafka クラスター用に、2 つ目の Kafka テーブルを作成します。
CREATE TABLE kafka_testing.second_kafka_table ON CLUSTER STAGE_CLICKHOUSE_CLUSTER
(
`id` UInt32,
`first_name` String,
`last_name` String
)
ENGINE = Kafka(cluster_2);レプリケートテーブルを作成する
最初のKafkaテーブル用に、次のテーブルを作成します:
CREATE TABLE kafka_testing.first_replicated_table ON CLUSTER STAGE_CLICKHOUSE_CLUSTER
(
`id` UInt32,
`first_name` String,
`last_name` String
) ENGINE = ReplicatedMergeTree()
ORDER BY id;2つ目のKafkaテーブル用のテーブルを作成します:
CREATE TABLE kafka_testing.second_replicated_table ON CLUSTER STAGE_CLICKHOUSE_CLUSTER
(
`id` UInt32,
`first_name` String,
`last_name` String
) ENGINE = ReplicatedMergeTree()
ORDER BY id;materialized view を作成する
最初の Kafka テーブルから最初のレプリケートテーブルへデータを挿入する materialized view を作成します。
CREATE MATERIALIZED VIEW kafka_testing.cluster_1_mv ON CLUSTER STAGE_CLICKHOUSE_CLUSTER TO first_replicated_table AS
SELECT
id,
first_name,
last_name
FROM first_kafka_table;2 つ目の Kafka テーブルから 2 つ目のレプリケートテーブルへデータを挿入するための materialized view を作成します。
CREATE MATERIALIZED VIEW kafka_testing.cluster_2_mv ON CLUSTER STAGE_CLICKHOUSE_CLUSTER TO second_replicated_table AS
SELECT
id,
first_name,
last_name
FROM second_kafka_table;セットアップの確認
これで、Kafkaクラスター上に対応するコンシューマグループが表示されているはずです。
cluster_1_clickhouse_consumer(cluster_1上)cluster_2_clickhouse_consumer(cluster_2上)
両方のテーブルのデータを確認するには、いずれかの ClickHouse ノードで次のクエリを実行します。
SELECT * FROM first_replicated_table LIMIT 10;SELECT * FROM second_replicated_table LIMIT 10;注
このガイドでは、両方のKafkaトピックに取り込まれるデータは同一です。実際の環境では、異なるデータになるはずです。Kafkaクラスターは必要な数だけ追加できます。
出力例:
┌─id─┬─first_name─┬─last_name─┐
│ 0 │ FirstName0 │ LastName0 │
│ 1 │ FirstName1 │ LastName1 │
│ 2 │ FirstName2 │ LastName2 │
└────┴────────────┴───────────┘これで、named collections を使用して ClickHouse と Kafka を統合するためのセットアップは完了です。Kafka の設定を ClickHouse の config.xml ファイルに集約することで、設定の管理や調整が容易になり、よりシンプルで効率的なインテグレーションを実現できます。