1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

Apache Iceberg Sink Connectorを使ってConfluent Cloudのトピックデータをwatsonx.dataへ連携する

1
Posted at

以下に記載されている「watsonx.dataへのConfluent Apache Iceberg Sink Connectorの統合」をやってみました!リンク先では S3 を利用した例が記載されていますが、本記事では ICOS を使って試してみます。またS3 の設定値についても併記しています。

Confluent の Apache Iceberg Sink Connectorは Kafka Connect のプラグインであり、 Kafka のトピックから取得したデータを、Exactly-Onceセマンティクスに基づいて Apache Iceberg のテーブルに直接書き込むことで、 watsonx.data からのストリーミングデータに対するリアルタイム分析を可能にします。

1. Iceberg Sinkコネクタ用の watsonx.data の設定

前提

  • Confluent Cloud からアクセス可能な S3-compatible ストレージバケット。ここ( AWS S3、 ICOS, MinIO,、または Ceph など)で、バケットのエンドポイント、アクセスキー、およびシークレットキーが利用可能なもの。 ここではICOSを使います
  • watsonx.data にプロビジョニングされた Presto エンジン。

1-1. S3 または ICOS ストレージを使用して、Icebergカタログを作成します。

  1. watsonx.data コンソールにログインする

  2. 左メニュー インフラストラクチャー・マネージャー を開く

  3. 右上のコンポーネントの追加 をクリック

  4. ストレージからICOSの場合は「IBM Cloud Object Storage」、 S3の場合は「Amazon S3」を選択し、「次へ」をクリック

  5. 必要な情報を入力後、「接続テスト」をクリック、テスト接続の正常完了を確認
    image.png

    image.png

  6. 「カタログの関連付け」をONにして、カタログ・タイプをApache Iceberg、カタログ名、基本パスを入力して「関連付け」をクリック

    image.png
    インフラストラクチャー・マネージャーの画面に戻ります。

1-2. Presto エンジンへのカタログの関連付け

  1. Presto のエンジンを探し、そこにマウスを重ねると「関連付けの管理」のアイコンが表示されますので、それをクリックします。
    image.png

  2. 1-1で作成したカタログ名にチェックを入れて、「保存してエンジンを再起動する」をクリックしてください。
    尚エンジンを再起動すると、クライアントからの接続は切断され、実施しているQueryは途中で終了してしまいますので、タイミングは注意してください。
    image.png

1-3. MDSエンドポイントを取得します。

Iceberg Sink Connectorが watsonx.data カタログサービスと通信するには、Metastore External REST (MDS) エンドポイントが必要です。

  1. 1-1で作成したカタログをクリック
    image.png

  2. メタストア REST エンドポイントをコピー
    Confluent Cloudでコネクタを設定する際は、メタストア REST エンドポイント(MDSエンドポイント)が必要です。 次の手順でこの値が利用できるようにしておいてください。
    image.png

2. Confluent Cloud で Iceberg Sink Connector を構成

今回は、Java のビルド作業が不要な、以下の公式ドキュメントの「Option 2: Download the pre-built ZIP from Confluent Hub (simpler, older version)」を利用します。

  1. 必要な Hadoop のJARファイルがすでに含まれている既製の iceberg-kafka-connect-1.9.2-fixed.zipをダウンロード

    https://www.ibm.com/docs/ja/SSAO5N/lh-over/topics/iceberg-kafka-connect-with-orc.zip

    ドキュメントにリンクのあるiceberg-kafka-connect-1.9.2-fixed.zipをダウンロードします。このファイルは2026/08/24現在上記サイトからダウンロード可能です。将来的にはバージョン等が変わる可能性があります。

    なお https://www.confluent.io/hub/iceberg/iceberg-kafka-connect からダウンロードしたZIPは Hadoop のJARファイルがなく、エラーがでて動きませんでした。

  2. Confluent Hub - Iceberg Connectorにアクセス
    以下にアクセスします。
    https://www.confluent.io/hub/iceberg/iceberg-kafka-connect

  3. 一番右のConfluent Cloudの「Launch on Cloud」をクリック
    元のダウンロードページ https://www.confluent.io/hub/iceberg/iceberg-kafka-connect の「Launch on Cloud」です。
    image.png

    以下のような画面が表示されたら「Continue」をクリック
    image.png

    Confluent Cloud にログインしていない場合はログイン画面が表示されるので、ログインします。
    image.png

  4. Add Custom Connector pluginの画面で詳細を入力
    Connector plugin details

    項目
    Connector plugin name 任意の名前 (ここではapache-iceberg-sink)
    Custom plugin description 任意 (ここではApache Iceberg Sink Connector for watsonx.data )
    Connector class org.apache.iceberg.connect.IcebergSinkConnector
    Connector type Sink

    Connector archive
    「Select connector archive」をクリックして、2-1でダウンロードしたZIPファイルをアップロードしてください。

    Sensitive properties

    Sensitive Property Key
    iceberg.catalog.rest.auth.basic.password
    iceberg.catalog.s3.access-key-id
    iceberg.catalog.s3.secret-access-key
    iceberg.kafka.sasl.jaas.config
  5. 最後に確認チェックボックスをオンにして 「Submit」をクリックします。

    image.png

3. Confluent Cloud にデータトピックを作成

Iceberg Sink Connector が Iceberg テーブルに書き込むストリーミングデータを受け取るトピックを作成します。

  1. 対象のConfluent Cloudクラスター画面を開きます

  2. 左のメニューから[Topics] クリック

  3. 「Create Topic」 をクリック。(すでにTopicが存在する場合は「Add Topic」)
    image.png

  4. topic name、Partitionsを入力
    ここではorder_topicと1としました

  5. [Create with default] をクリック
    image.png

    以下の画面が表示されたら、「Skip」をクリックしてください
    image.png

4. Confluent Cloud にコントロール・トピックを作成

コントロール・トピックは、Iceberg Sink Connector のすべてのタスクにわたるコミットを調整し、「正確に1回」のセマンティクスを保証します。 複数のコネクタタスクが同時に同じIcebergテーブルに書き込みを行う際、データの破損を防ぐことができます。

  1. 対象のConfluent Cloudクラスター画面を開きます

  2. 左のメニューから[Topics] クリック

  3. 「Add Topic」 をクリック。
    image.png

  4. topic name、Partitionsに1を入力
    ここではcontrol-icebergと1としました

  5. [Create with default] をクリック

    image.png

    以下の画面が表示されたら、「Skip」をクリックしてください
    image.png

5. 上流のデータソースを作成

テストを行うには、Datagen Source Connector を使用してサンプルデータを生成し、それを Kafka トピックに公開してください。 実際のデータソースがある場合は、この手順をスキップしてください。

  1. 対象のConfluent Cloudクラスター画面を開きます

  2. 左のメニューから「Connectors」をクリック

  3. 「Sample Data」のタイルをクリックします。
    image.png

  4. 表示されたウィンドウで「Additional Configurationn」をクリックします。
    image.png

  5. 3 で作成したデータトピック(ここではorder_topic)を選択し、「Continue」 をクリック
    image.png

  6. 「Generate API Key and download」 をクリックし、API Keyを作成して保管します。
    image.png

  7. 「Continue」 をクリックしてください。
    image.png

  8. 出力形式を設定します

    • Select output record value format: JSONを選択
    • Select a schema:ordersを選択
      「Continue」 をクリックし、デフォルトのサイズ設定を受け入れて、もう一度「Continue」をクリックします。
      image.png
      image.png
  9. 設定を確認し、「Continue」 をクリックしてコネクタを起動します
    image.png

6. Iceberg Sink Connectorを作成

  1. 対象のConfluent Cloudクラスター画面を開きます

  2. 左のメニューから「Connectors」をクリック

  3. 「See all connectors」をクリック後、2.4で設定した名前で検索(ここではapache-iceberg-sink)し、そのタイルをクリック

    image.png
    image.png

  4. 「Use an existing API key」 をクリックし、5の6で作成したAPI Keyと
    Secretを入力して「Continue」をクリックします。
    image.png

  5. Configurationで「 JSON 」タブをクリックし、以下のコネクタ設定を入力します。
    image.png

    {
    "iceberg.catalog": "new_catalog",
    "iceberg.catalog.client.region": "us-west-2",
    "iceberg.catalog.header.AccountId": "<account id>",
    "iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
    "iceberg.catalog.rest.auth.basic.password": "<API key>",
    "iceberg.catalog.rest.auth.basic.username": "<username>",
    "iceberg.catalog.rest.auth.type": "basic",
    "iceberg.catalog.s3.access-key-id": "<Access key>",
    "iceberg.catalog.s3.endpoint": "https://s3.<region>.amazonaws.com",
    "iceberg.catalog.s3.path-style-access": "true",
    "iceberg.catalog.s3.secret-access-key": "<secret key>",
    "iceberg.catalog.type": "rest",
    "iceberg.catalog.uri": "<MDS rest endpoint>/api/v1/iceberg",
    "iceberg.catalog.warehouse": "<watsonx.data catalog associated with S3 bucket>",
    "iceberg.control.commit.interval-ms": "300000",
    "iceberg.control.commit.timeout-ms": "30000",
    "iceberg.control.topic": "<control-topic>",
    "iceberg.kafka.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"<Confluent cluster API key>\"    password=\"<Confluent cluster secret>\";",
    "iceberg.kafka.sasl.mechanism": "PLAIN",
    "iceberg.kafka.security.protocol": "SASL_SSL",
    "iceberg.tables": "default.<topic>",
    "iceberg.tables.auto-create-enabled": "true",
    "iceberg.tables.default-compression-codec": "snappy",
    "iceberg.tables.default-file-format": "parquet",
    "iceberg.tables.evolve-schema-enabled": "true",
    "iceberg.tables.schema-force-optional": "false",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "key.converter.schemas.enable": "false",
    "topics": "<topic>",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false"
    }
    

    上記の<>で囲まれた値を自分の環境の値に変更します:

    設定項目 説明
    iceberg.catalog.header.AccountId watsonx.dataを使用するアカウントの Account ID IBM Cloudのwebコンソール 管理アカウントアカウント設定から取得
    iceberg.catalog.rest.auth.basic.password IBM Cloud API Key watsonx.data の権限があるIDのAPI Key
    iceberg.catalog.rest.auth.basic.username 上記のiceberg.catalog.rest.auth.basic.passwordの権限があるAPI KeyのIDをibmlhapikey_<ID>の形式にする ibmlhapikey_xxxx@ibm.com など
    iceberg.catalog.s3.access-key-id S3 Access Key S3: IAM ユーザーのアクセスキー
    ICOS: ICOS HMAC access_key_id
    iceberg.catalog.s3.endpoint S3 エンドポイント S3: S3エンドポイントを指定
    ICOS: ICOSエンドポイント https://s3.<region>.cloud-object-storage.appdomain.cloud
    iceberg.catalog.s3.secret-access-key S3 Secret Key S3: IAM ユーザーのシークレットキー
    ICOS: HMAC secret_access_key
    iceberg.catalog.uri <メタストア REST エンドポイント>/api/v1/iceberg 1-3で取得したメタストア REST エンドポイント
    iceberg.catalog.warehouse watsonx.dataのカタログ名 1-1の6で設定したもの
    "iceberg.control.topic コントロール・トピック名 4の4で設定したもの
    "iceberg.kafka.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"<Confluent cluster API key>\" password=\"<Confluent cluster secret>\";" <>にkafkaのAPI Keyのusernameとpasswordを入れる kafka権限があればいいので「Generate API Key and download」でダウンロードしたものどれかを入れる
    iceberg.tables . watsonx.dataでのスキーマ名とテーブル名
    topics 連携するトピック名 3の4で設定したもの
    Sample
    {
    "iceberg.catalog": "new_catalog",
    "iceberg.catalog.client.region": "us-west-2",
    "iceberg.catalog.header.AccountId": "0bbd999999d09fb9b9d9ab99c9a999e9",
    "iceberg.catalog.io-impl": "org.apache.iceberg.aws.s3.S3FileIO",
    "iceberg.catalog.rest.auth.basic.password": "Xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx",
    "iceberg.catalog.rest.auth.basic.username": "ibmlhapikey_xxxx@bm.com",
    "iceberg.catalog.rest.auth.type": "basic",
    "iceberg.catalog.s3.access-key-id": "99999999999999999999999999999999",
    "iceberg.catalog.s3.endpoint": "https://s3.ca-tor.cloud-object-storage.appdomain.cloud",
    "iceberg.catalog.s3.path-style-access": "true",
    "iceberg.catalog.s3.secret-access-key": "999999999999999999999999999999999999999999999999",
    "iceberg.catalog.type": "rest",
    "iceberg.catalog.uri": "https://console-ibm-cator.lakehouse.saas.ibm.com:443/api/v1/iceberg",
    "iceberg.catalog.warehouse": "confluent",
    "iceberg.control.commit.interval-ms": "300000",
    "iceberg.control.commit.timeout-ms": "30000",
    "iceberg.control.topic": "control-iceberg",
    "iceberg.kafka.sasl.jaas.config": "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"XXXXXXXXXXXXXXXX\"    password=\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\";",
    "iceberg.kafka.sasl.mechanism": "PLAIN",
    "iceberg.kafka.security.protocol": "SASL_SSL",
    "iceberg.tables": "demo_c.order_topic",
    "iceberg.tables.auto-create-enabled": "true",
    "iceberg.tables.default-compression-codec": "snappy",
    "iceberg.tables.default-file-format": "parquet",
    "iceberg.tables.evolve-schema-enabled": "true",
    "iceberg.tables.schema-force-optional": "false",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "key.converter.schemas.enable": "false",
    "topics": "order_topic",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false"
    }
    

    最後に「Continue」をクリックします。
    image.png

  6. Networkingで以下のエンドポイントを追加
    Confluent Cloud の Custom Connector はデフォルトで外部への通信がブロックされるため、以下のエンドポイントを追加します。

    エンドポイント 説明
    <メタストア REST エンドポイント> REST Catalog への接続
    S3 エンドポイント AWS S3 または ICOSへの書き込み

    例(ICOS & watsonx.data Toronto リージョン):

    console-ibm-cator.lakehouse.saas.ibm.com:443
    s3.ca-tor.cloud-object-storage.appdomain.cloud:443
    

    最後に「Continue」をクリックします。
    image.png

  7. Connector sizingはデフォルトのまま、「Continue」をクリックします。
    名称未設定15-1(ドラッグされました).tiff

  8. 設定を確認し、「Continue」をクリックします。
    image.png

  9. Connector画面で「Iceberg Sink Connector」が 「Running」 状態になっていることを確認してください。

image.png

これで設定終了です。

7. Confluent の「Iceberg Sink Connector」統合の検証

  1. S3またはICOS バケット内にデータファイルとメタデータファイルが作成されていることを確認します

    コネクタの起動後、 設定したバケット内の、設定 iceberg.tables 値に対応するネームスペースの下に、以下のディレクトリ構造が存在するか確認してください:

    • data/ - コネクタによって書き込まれたParquetデータファイルが格納されています。
    • metadata/ - マニフェストやスナップショットファイルなど、Icebergのメタデータファイルが含まれています。

    両方のディレクトリが存在していることから、コネクタがIcebergテーブルへのデータ書き込みに成功していることが確認できます。

    image.png

  2. watsonx.data 内のスキーマとテーブルを確認します
    watsonx.data コンソールで、「照会ワークスペースに移動し、エンジンのドロップダウンメニューから、お使いの Presto エンジンを選択します。
    以下のクエリを実行し、データが存在することを確認してください:

    SELECT count(*) FROM 
    "confluent"."demo_c"."order_topic";
    

    "confluent"."demo_c"."order_topic"は自分で作成した名前に変更してください。
    データ件数が表示されることを確認できました。
    image.png

watsonx.data のIcebergテーブルを通じて、 Kafka のトピックからリアルタイムのストリーミングデータをクエリできるようになりました。 Kafka トピックに公開された新しいメッセージは、コネクタのコミット間隔が経過すると、自動的にIcebergテーブルに反映されます。

まとめ

Apache Iceberg Sink ConnectorはOSSのコネクタですが、公式ドキュメントに従って、Confluentのトピックを watsonx.dataにデータ連携することが可能です。ぜひお試しください。
特に今回は Confluent Cloud を利用しましたが、オンプレミス版の Confluent Platform では Tableflow を利用できないため、watsonx.data との連携手段として Apache Iceberg Sink Connector が有力な選択肢になると思います。

(残念ながら環境が入手できなかったためオンプレミス版の Confluent Platform での検証は実施していません)

1
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?