組み込みの Kafka テーブルエンジンと、マテリアライズドビューを使用して、ApsaraMQ for Kafka から ApsaraDB for ClickHouse へデータをリアルタイムで同期できます。
制限事項
データは、ApsaraMQ for Kafka インスタンス、および ECS インスタンスにデプロイされたセルフマネージド Kafka クラスターからのみ同期できます。
前提条件
-
ApsaraDB for ClickHouse :
-
ApsaraMQ for Kafka インスタンスと同じリージョンおよび VPC 内にターゲットクラスターを作成しておく必要があります。詳細については、「クラスターの作成」をご参照ください。
-
ターゲットクラスター用に、必要な権限を持つデータベースアカウントを作成しておく必要があります。詳細については、「アカウント管理」をご参照ください。
-
-
ApsaraMQ for Kafka :
-
トピックを作成しておく必要があります。
-
コンシューマーグループを作成しておく必要があります。
-
注意事項
-
ApsaraDB for ClickHouse の Kafka 外部テーブルがサブスクライブするトピックには、他のコンシューマーが存在しない必要があります。
-
Kafka 外部テーブル、マテリアライズドビュー、ローカルテーブルを作成する際は、3 つのテーブルのフィールドの型が一致している必要があります。
手順
次の例では、ApsaraMQ for Kafka から ApsaraDB for ClickHouse Community-compatible Edition クラスターの default データベースにある kafka_table_distributed 分散テーブルにデータを同期します。
手順1:同期の仕組みの理解
ApsaraDB for ClickHouse は、Kafka テーブルエンジンとマテリアライズドビューを使用して、Kafka からリアルタイムでデータを消費および格納します。 データフローは次のとおりです:
-
Kafka トピック:同期対象のソースデータ。
-
ApsaraDB for ClickHouse の Kafka 外部テーブル (Kafka テーブルエンジンを使用するテーブル) :指定された Kafka トピックからソースデータをプルします。
-
マテリアライズドビュー: Kafka の外部テーブルからソースデータを読み取り、ApsaraDB for ClickHouse のローカルテーブルにデータを挿入します。
-
ローカルテーブル:同期されたデータを格納します。
手順2: ApsaraDB for ClickHouse クラスターへの接続
詳細については、「DMS を使用した ApsaraDB for ClickHouse クラスターへの接続」をご参照ください。
手順3:Kafka 外部テーブルの作成
Kafka 外部テーブルは、Kafka テーブルエンジンを使用して、指定された Kafka トピックからデータをプルします。 このテーブルには次の特徴があります。
-
デフォルトでは、Kafka 外部テーブルを直接クエリすることはできません。
-
Kafka 外部テーブルは Kafka データの消費にのみ使用され、データは格納しません。 マテリアライズドビューを使用してデータを処理し、宛先テーブルに挿入する必要があります。
テーブルを作成するための構文は次のとおりです:
Kafka 外部テーブルのフィールドの型は、Kafka 内のメッセージのデータ型と一致している必要があります。
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'host:port1,host:port2,host:port3',
kafka_topic_list = 'topic_name1,topic_name2,...',
kafka_group_name = 'group_name',
kafka_format = 'data_format'[,]
[kafka_row_delimiter = 'delimiter_symbol',]
[kafka_num_consumers = N,]
[kafka_thread_per_consumer = 1,]
[kafka_max_block_size = 0,]
[kafka_skip_broken_messages = N,]
[kafka_commit_every_batch = 0,]
[kafka_auto_offset_reset = 'value']
一般的なパラメーターを次の表で説明します:
|
パラメーター |
必須 |
説明 |
|
kafka_broker_list |
はい |
Kafka クラスターのブローカーエンドポイントのカンマ区切りリスト。 エンドポイントの表示方法の詳細については、「エンドポイントの表示」をご参照ください。
|
|
kafka_topic_list |
はい |
トピック名のカンマ区切りリスト。 トピック名の表示方法の詳細については、「トピックの作成」をご参照ください。 |
|
kafka_group_name |
はい |
Kafka コンシューマーグループの名前。 詳細については、「グループの作成」をご参照ください。 |
|
kafka_format |
はい |
ApsaraDB for ClickHouse が処理できるメッセージボディの形式です。 説明
ApsaraDB for ClickHouse がサポートするメッセージボディ形式の詳細については、「入出力データ形式」をご参照ください。 |
|
kafka_row_delimiter |
いいえ |
行を区切るための区切り文字です。 デフォルト値は \n です。 実際のデータで使用されている区切り文字に合わせてこのパラメーターを設定することもできます。 |
|
kafka_num_consumers |
いいえ |
1 つのテーブルあたりのコンシューマー数です。 デフォルト値は 1 です。 説明
|
|
kafka_thread_per_consumer |
いいえ |
各コンシューマーに専用のスレッドを有効にするかどうかを指定します。 デフォルト値は 0 です。 有効な値:
消費速度を向上させる方法の詳細については、「Kafka パフォーマンスチューニング」をご参照ください。 |
|
kafka_max_block_size |
いいえ |
Kafka メッセージバッチの最大サイズ (バイト単位) です。 デフォルト値は 65536 です。 |
|
kafka_skip_broken_messages |
いいえ |
無視する解析エラーの数です。 デフォルト値は 0 です。 |
|
kafka_commit_every_batch |
いいえ |
Kafka のコミット頻度です。 デフォルト値は 0 です。 有効な値:
|
|
kafka_auto_offset_reset |
いいえ |
Kafka データの読み取りを開始するオフセットです。 有効な値:
説明
このパラメーターは、カーネルバージョンが 21.8 の ApsaraDB for ClickHouse クラスターではサポートされていません。 |
パラメーターの詳細については、「Kafka」をご参照ください。
次のコードに例を示します:
CREATE TABLE default.kafka_src_table ON CLUSTER `default`
(
-- テーブルスキーマのフィールドを定義します。
id Int32,
name String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****1-2-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****-3-vpc.alikafka.aliyuncs.com:9092',
kafka_topic_list = 'testforCK',
kafka_group_name = 'GroupForTestCK',
kafka_format = 'CSV';
手順4:宛先テーブルの作成
クラスターのエディションに応じて、テーブル作成ステートメントを選択します。
Enterprise Edition クラスターの場合は、ローカルテーブルを作成するだけで済みます。 Community-compatible Edition クラスターの場合は、環境と要件に基づいて分散テーブルを作成する必要がある場合があります。 以下にステートメントの例を示します。 テーブル作成の構文の詳細については、「CREATE TABLE」をご参照ください。
Enterprise Edition
CREATE TABLE default.kafka_table_local ON CLUSTER default (
id Int32,
name String
) ENGINE = MergeTree()
ORDER BY (id);
このステートメントを実行したときに ON CLUSTER is not allowed for Replicated database エラーが表示された場合は、カーネルバージョンをアップグレードすることで問題を解決できます。カーネルバージョンのアップグレード方法の詳細については、「マイナーエンジンバージョンのアップグレード」をご参照ください。
Community-compatible Edition
シングルレプリカクラスターとダブルレプリカクラスターではテーブルエンジンが異なります。 クラスターのレプリカタイプに基づいて適切なエンジンを選択してください。
ダブルレプリカクラスターでテーブルを作成する場合、MergeTree エンジンファミリーの Replicated エンジンを使用する必要があります。ダブルレプリカクラスターで Replicated 以外のエンジンを使用してテーブルを作成すると、レプリカ間でデータがレプリケーションされず、データの不整合が発生する可能性があります。
シングルレプリカ
-
ローカルテーブルを作成します。
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = MergeTree() ORDER BY (id); -
(オプション) 分散テーブルを作成します。
データをローカルテーブルにインポートするだけでよい場合は、この手順をスキップしてください。
マルチノードクラスターを使用している場合は、分散テーブルを作成することを推奨します。
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
ダブルレプリカ
-
ローカルテーブルを作成します。
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = ReplicatedMergeTree() ORDER BY (id); -
(オプション) 分散テーブルを作成します。
データをローカルテーブルにインポートするだけでよい場合は、この手順をスキップしてください。
マルチノードクラスターを使用している場合は、分散テーブルを作成することを推奨します。
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
手順5:マテリアライズドビューの作成
ApsaraDB for ClickHouse は、マテリアライズドビューを使用して Kafka 外部テーブルからソースデータを読み取り、ApsaraDB for ClickHouse のローカルテーブルにデータを挿入します。
マテリアライズドビューを作成するための構文は次のとおりです:
SELECT フィールドが宛先テーブルの構造と一致することを確認するか、変換関数を使用してデータ形式を宛先テーブルの構造と一致させてください。
CREATE MATERIALIZED VIEW <view_name> ON CLUSTER default TO <dest_table> AS SELECT * FROM <src_table>;
次の表に、パラメーターを示します。
|
パラメーター |
必須 |
説明 |
例 |
|
view_name |
はい |
ビューの名前です。 |
consumer |
|
dest_table |
はい |
Kafka データを格納する宛先テーブルです。
|
|
|
src_table |
はい |
Kafka 外部テーブルです。 |
kafka_src_table |
以下にステートメントの例を示します:
Enterprise Edition
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_local AS SELECT * FROM kafka_src_table;
Community-compatible Edition
この例では、ソースデータは kafka_table_distributed 分散テーブルに格納されます。
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_distributed AS SELECT * FROM kafka_src_table;
ステップ 6:同期の検証
-
ApsaraMQ for Kafka インスタンスのトピックにメッセージを送信します。
-
ApsaraMQ for Kafka コンソールにログインします。
-
インスタンスリスト ページで、対象インスタンスの名前をクリックします。
-
[トピック] ページで、対象トピックを見つけ、操作 列で を選択します。
-
[クイック体験によるメッセージの送受信] ページで、Message Content を入力します。
この例では、
1,aおよび2,bのメッセージを送信します。 -
を決定 をクリックします。
-
-
ApsaraDB for ClickHouse クラスターにログオンし、分散テーブルをクエリして、データが同期されているかどうかを確認します。
ApsaraDB for ClickHouse クラスターへのログオン方法の詳細については、「DMS を使用して ApsaraDB for ClickHouse クラスターに接続する」をご参照ください。
次のステートメントを使用してデータをクエリし、検証します。
Enterprise Edition
SELECT * FROM kafka_table_local;Community-compatible Edition
次のコードは、分散テーブルをクエリする例を示しています。
-
宛先テーブルがローカルテーブルの場合は、クエリ内の 分散テーブル 名をローカルテーブル名に置き換える必要があります。
-
マルチノードクラスターである Community-compatible Edition クラスター を使用する場合は、分散テーブルをクエリすることを 強く推奨します。ローカルテーブルを直接クエリすると、単一ノードのデータのみが返され、結果セットが不完全になります。
SELECT * FROM kafka_table_distributed;クエリで結果が返されると、Kafka から ApsaraDB for ClickHouse へのデータ同期は成功です。
クエリ結果は次のとおりです。
┌─id─┬─name─┐ │ 1 │ a │ │ 2 │ b │ └────┴──────┘クエリ結果が期待どおりでない場合は、「ステップ 7 (オプション):Kafka 外部テーブルの消費ステータスの確認」に進み、問題のトラブルシューティングを行ってください。
-
ステップ 7 (オプション):Kafka 消費ステータスの確認
同期されたデータが Kafka 内のデータと一致しない場合は、システムテーブルをクエリして Kafka 外部テーブルの消費ステータスを確認し、例外のトラブルシューティングを行います。
エンジン v23.8 以降
次のステートメントを実行して system.kafka_consumers システムテーブルをクエリし、Kafka 外部テーブルの消費ステータスを表示します。
select * from system.kafka_consumers;
次の表は、system.kafka_consumers テーブルのフィールドについて説明しています。
|
フィールド |
説明 |
|
database |
Kafka 外部テーブルが配置されているデータベース。 |
|
table |
Kafka 外部テーブルの名前。 |
|
consumer_id |
Kafka コンシューマーの ID。 1 つのテーブルに複数のコンシューマーを設定できます。コンシューマー数は、Kafka 外部テーブル作成時に kafka_num_consumers パラメータで指定します。 |
|
assignments.topic |
Kafka トピック。 |
|
assignments.partition_id |
Kafka パーティションの ID。 1 つのパーティションに割り当てることができるコンシューマーは 1 つのみです。 |
|
assignments.current_offset |
現在のオフセット。 |
|
exceptions.time |
直近 10 件の例外のタイムスタンプ。 |
|
exceptions.text |
直近 10 件の例外のテキスト。 |
|
last_poll_time |
最後のポーリングのタイムスタンプ。 |
|
num_messages_read |
コンシューマーが読み取ったメッセージ数。 |
|
last_commit_time |
最後のコミットのタイムスタンプ。 |
|
num_commits |
コンシューマーが実行したコミットの総数。 |
|
last_rebalance_time |
最後の Kafka リバランスのタイムスタンプ。 |
|
num_rebalance_revocations |
コンシューマーからパーティションが取り消された回数。 |
|
num_rebalance_assignments |
Kafka クラスター内でコンシューマーにパーティションが割り当てられた回数。 |
|
is_currently_used |
コンシューマーが使用中かどうかを示します。 |
|
last_used |
コンシューマーが最後に使用された時刻 (UNIX 時間、マイクロ秒単位)。 |
|
rdkafka_stat |
ライブラリの内部統計情報。詳細については、「librdkafka」をご参照ください。 デフォルト値は 3000 で、3 秒ごとに統計情報が生成されることを示します。 説明
ApsaraDB for ClickHouse で |
エンジン v23.8 より前
次のステートメントを実行して system.kafka システムテーブルをクエリし、Kafka 外部テーブルの消費ステータスを表示します。
SELECT * FROM system.kafka;
次の表は、system.kafka テーブルのフィールドについて説明しています。
|
フィールド |
説明 |
|
database |
Kafka 外部テーブルが配置されているデータベースの名前。 |
|
table |
Kafka 外部テーブルの名前。 |
|
topic |
Kafka 外部テーブルが消費するトピックの名前。 |
|
consumer_group |
Kafka 外部テーブルが使用するコンシューマーグループの名前。 |
|
last_read_message_count |
Kafka 外部テーブルからプルされたメッセージ数。 |
|
status |
外部テーブルによる Kafka メッセージ消費のステータス。有効な値:
|
|
exception |
例外の詳細。 説明
status の値が error の場合、このパラメータは例外の詳細を返します。 |