すべてのプロダクト
Search
ドキュメントセンター

ApsaraDB for ClickHouse:Kafka からのデータ同期

最終更新日:Aug 27, 2026

組み込みの 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 からリアルタイムでデータを消費および格納します。 データフローは次のとおりです:

image
  • 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 クラスターのブローカーエンドポイントのカンマ区切りリスト。 エンドポイントの表示方法の詳細については、「エンドポイントの表示」をご参照ください。

  • ApsaraMQ for Kafka を使用する場合、ApsaraDB for ClickHouse はデフォルトで ApsaraMQ for Kafka インスタンスのドメイン名を解析できます。

  • 自己管理型 Kafka クラスターを使用する場合、ApsaraDB for ClickHouse では、IP アドレスまたは固定形式のカスタムドメイン名を使用して Kafka クラスターに接続できます。 次のカスタムドメイン名ルールをサポートしています。

    1. .com で終わるドメイン名。

    2. .local で終わり、kafka、mysql、または rabbitmq を含むドメイン名。

kafka_topic_list

はい

トピック名のカンマ区切りリスト。 トピック名の表示方法の詳細については、「トピックの作成」をご参照ください。

kafka_group_name

はい

Kafka コンシューマーグループの名前。 詳細については、「グループの作成」をご参照ください。

kafka_format

はい

ApsaraDB for ClickHouse が処理できるメッセージボディの形式です。

説明

ApsaraDB for ClickHouse がサポートするメッセージボディ形式の詳細については、「入出力データ形式」をご参照ください。

kafka_row_delimiter

いいえ

行を区切るための区切り文字です。 デフォルト値は \n です。 実際のデータで使用されている区切り文字に合わせてこのパラメーターを設定することもできます。

kafka_num_consumers

いいえ

1 つのテーブルあたりのコンシューマー数です。 デフォルト値は 1 です。

説明
  1. 1 つのコンシューマーのスループットが不十分な場合は、コンシューマーを増やす必要があります。

  2. 各パーティションに割り当てることができるコンシューマーは 1 つだけであるため、コンシューマーの総数はトピック内のパーティションの数を超えることはできません。

kafka_thread_per_consumer

いいえ

各コンシューマーに専用のスレッドを有効にするかどうかを指定します。 デフォルト値は 0 です。 有効な値:

  1. 0:すべてのコンシューマーが 1 つのスレッドを共有してデータを消費します。

  2. 1:各コンシューマーが専用のスレッドを有効にしてデータを消費します。

消費速度を向上させる方法の詳細については、「Kafka パフォーマンスチューニング」をご参照ください。

kafka_max_block_size

いいえ

Kafka メッセージバッチの最大サイズ (バイト単位) です。 デフォルト値は 65536 です。

kafka_skip_broken_messages

いいえ

無視する解析エラーの数です。 デフォルト値は 0 です。 kafka_skip_broken_messages=N を設定すると、エンジンは N 個の解析できない Kafka メッセージをスキップします。 1 つのメッセージは 1 行のデータに相当します。

kafka_commit_every_batch

いいえ

Kafka のコミット頻度です。 デフォルト値は 0 です。 有効な値:

  1. 0:完全なデータブロックが書き込まれた後にのみコミットが実行されます。

  2. 1:データの各バッチが書き込まれた後にコミットが実行されます。

kafka_auto_offset_reset

いいえ

Kafka データの読み取りを開始するオフセットです。 有効な値:

  1. earliest:最も古いオフセットから Kafka データを読み取ります。 これがデフォルト値です。

  2. latest:最新のオフセットから 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 以外のエンジンを使用してテーブルを作成すると、レプリカ間でデータがレプリケーションされず、データの不整合が発生する可能性があります。

シングルレプリカ

  1. ローカルテーブルを作成します。

    CREATE TABLE default.kafka_table_local ON CLUSTER default (
      id Int32,
      name String
    ) ENGINE = MergeTree()
    ORDER BY (id);
  2. (オプション) 分散テーブルを作成します。

    データをローカルテーブルにインポートするだけでよい場合は、この手順をスキップしてください。

    マルチノードクラスターを使用している場合は、分散テーブルを作成することを推奨します。

    CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local
    ENGINE = Distributed(default, default, kafka_table_local, id);

ダブルレプリカ

  1. ローカルテーブルを作成します。

    CREATE TABLE default.kafka_table_local ON CLUSTER default (
      id Int32,
      name String
    ) ENGINE = ReplicatedMergeTree()
    ORDER BY (id);
  2. (オプション) 分散テーブルを作成します。

    データをローカルテーブルにインポートするだけでよい場合は、この手順をスキップしてください。

    マルチノードクラスターを使用している場合は、分散テーブルを作成することを推奨します。

    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 データを格納する宛先テーブルです。

  • Community-compatible Edition クラスター:

    • マルチノードクラスターの場合は、分散テーブルにデータをインポートすることを推奨します。

    • ローカルテーブルに同期する場合は、ローカルテーブル名を指定してください。

  • Enterprise Edition クラスター: Enterprise Edition クラスターには分散テーブルがないため、ローカルテーブルを指定する必要があります。

  • Community-compatible Edition の例:kafka_table_distributed

  • Enterprise Edition の例:kafka_table_local

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:同期の検証

  1. ApsaraMQ for Kafka インスタンスのトピックにメッセージを送信します。

    1. ApsaraMQ for Kafka コンソールにログインします。

    2. インスタンスリスト ページで、対象インスタンスの名前をクリックします。

    3. [トピック] ページで、対象トピックを見つけ、操作 列で より > [メッセージ送信 (デモ)] を選択します。

    4. [クイック体験によるメッセージの送受信] ページで、Message Content を入力します。

      この例では、1,a および 2,b のメッセージを送信します。

    5. を決定 をクリックします。

  2. 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 で statistics_interval_ms=0 が設定されている場合、Kafka 外部テーブルの統計収集は無効になります。

エンジン 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 メッセージ消費のステータス。有効な値:

  • no_view:Kafka 外部テーブルに対してビューが作成されていません。

  • attach_view:Kafka 外部テーブルに対してビューが作成されています。

  • 正常:ステータスは正常です。

    正常ステータスは、外部テーブルが期待どおりにデータを消費していることを示します。

  • skip_parse:解析エラーがスキップされています。

  • error:消費例外が発生しました。

exception

例外の詳細。

説明

status の値が error の場合、このパラメータは例外の詳細を返します。

よくある質問