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

Data Lake Formation:Iceberg REST を使用した Flink CDC からの DLF カタログへのアクセス

最終更新日:Jun 29, 2026

Flink CDC を使用して、Realtime Compute for Apache Flink の Iceberg REST を介してデータを Data Lake Formation (DLF) カタログに同期します。

前提条件

  • Realtime Compute for Apache Flink のフルマネージドワークスペースが必要です。作成していない場合は、「Realtime Compute for Apache Flink の有効化」をご参照ください。

  • Realtime Compute for Apache Flink ワークスペースと DLF が同じリージョンにあることを確認してください。また、ワークスペースの VPC を DLF ホワイトリストに追加する必要があります。詳細については、「VPC ホワイトリストの設定」をご参照ください。

制限事項

DLF への Iceberg REST 接続には、Realtime Compute for Apache Flink エンジンバージョン VVR 11.6.0 以降が必要です。

Flink での DLF カタログの登録

説明
  • この操作により、DLF カタログへのマッピングが作成されます。Flink でカタログを作成または削除しても、DLF 内の実際のデータには影響しません。

  • Iceberg REST を介して DLF カタログに作成されたすべてのテーブルは、Iceberg テーブルです。

  1. Realtime Compute for Apache Flink 管理コンソールにログインします。

  2. ワークスペースの [アクション] 列で、[コンソール] をクリックします。

  3. 左側のナビゲーションメニューで、[開発] > [スクリプト] をクリックします。

  4. 新しいスクリプトを作成します。SQL エディターで、次の SQL 文をコピーして貼り付けます。右下隅で [環境] をクリックし、VVR 11.2.0 以降のセッションクラスターを選択して、SQL 文を実行し、Iceberg REST を介して DLF カタログを登録します。

    CREATE CATALOG `catalog_name` 
     WITH (
        'type' = 'iceberg',
        'catalog-type' = 'rest',
        'uri' = 'http://cn-hangzhou-vpc.dlf.aliyuncs.com/iceberg',
        'warehouse' = 'iceberg_test',
        'rest.signing-region' = 'cn-hangzhou',
        'io-impl' = 'org.apache.iceberg.rest.DlfFileIO'
    );

    次の表にオプションを説明します。

    オプション

    説明

    必須

    type

    タイプ。これを iceberg に設定します。

    はい

    iceberg

    catalog-type

    カタログタイプ。これを rest に設定します。

    はい

    rest

    token.provider

    トークンプロバイダー。これを dlf に設定します。

    はい

    dlf

    uri

    Iceberg REST を介して DLF カタログにアクセスするために使用される URI。詳細については、「Iceberg REST」をご参照ください。

    はい

    http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg

    warehouse

    DLF カタログの名前。

    はい

    iceberg_test

    rest.signing-region

    DLF のリージョン ID。詳細については、「エンドポイント」をご参照ください。

    はい

    ap-southeast-1

    io-impl

    これを org.apache.iceberg.rest.DlfFileIO に設定します。

    はい

    org.apache.iceberg.rest.DlfFileIO

カタログを使用するための Flink CDC の設定

データインジェストジョブを作成するには、「Flink CDC データインジェストジョブの開発」をご参照ください。

すでに Flink カタログマッピングを作成している場合は、「既存のカタログを再利用して接続情報を取得する」を参照し、データインジェストジョブに次のシンク設定を追加します。

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

設定例

以下の例は、Flink CDC YAML ジョブを使用した一般的なデータ同期パターンを示しています。

MySQL データベース全体の DLF への同期

このジョブは、MySQL データベース全体を DLF に同期します。

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (オプション) 増分フェーズ中に作成されたテーブルを同期します。
  scan.binlog.newly-added-table.enabled: true
  # (オプション) テーブルとフィールドのコメントを同期します。
  include-comments.enabled: true
  # (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name
説明

MySQL ソースでは、これらのオプションパラメーターの使用を推奨します。詳細については、「MySQL」をご参照ください。

  1. パラメーター:scan.binlog.newly-added-table.enabled

    機能:増分フェーズ中に作成されたテーブルを同期します。

  2. パラメーター:include-comments.enabled

    機能:テーブルとフィールドのコメントを同期します。

  3. パラメーター:scan.incremental.snapshot.unbounded-chunk-first.enabled

    機能:TaskManager の OOM エラーを防ぎます。

  4. パラメーター:scan.only.deserialize.captured.tables.changelog.enabled

    機能:一致したテーブルのみを解析することで読み取りを高速化します。

パーティション分割された DLF テーブルへの書き込み

partition-keys パラメーターを使用してパーティションキーを指定します。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。例:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (オプション) 増分フェーズ中に作成されたテーブルを同期します。
  scan.binlog.newly-added-table.enabled: true
  # (オプション) テーブルとフィールドのコメントを同期します。
  include-comments.enabled: true
  # (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (オプション) パーティションキーを設定します。
    partition-keys: id,pt
  - source-table: mysql_test.tbl2
    partition-keys: id,pt

追記専用の DLF テーブルへの書き込み

データインジェストジョブは、ソースからのすべての変更イベントタイプをキャプチャします。DELETE 操作をダウンストリームの INSERT 操作に変換してソフトデリートを実装するには、以下のようにジョブを設定します。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (オプション) 増分フェーズ中に作成されたテーブルを同期します。
  scan.binlog.newly-added-table.enabled: true
  # (オプション) テーブルとフィールドのコメントを同期します。
  include-comments.enabled: true
  # (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (オプション) パーティションキーを設定します。
    partition-keys: id,pt
    # (オプション) ソフトデリートを実装します。
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
  - source-table: mysql_test.tbl2
    # (オプション) パーティションキーを設定します。
    partition-keys: id,pt
    # (オプション) ソフトデリートを実装します。
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
説明
  • プロジェクションに __data_event_type__ を追加すると、変更イベントタイプがダウンストリームテーブルに新しいフィールドとして書き込まれます。converter-after-transformSOFT_DELETE に設定すると、DELETE 操作が INSERT に変換されます。これにより、ダウンストリームテーブルにはすべての変更イベントが記録されます。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。

リアルタイム Kafka CDC データの DLF への同期

Kafka の inventory トピックに、customersproducts の 2 つのテーブルの変更データが Debezium JSON 形式で保存されているとします。このジョブは、これらのテーブルのデータを対応する DLF ターゲットに同期します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: debezium-json
  debezium-json.distributed-tables: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Debezium JSON には主キー情報がないため、手動で追加します。
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
説明
  • Kafka ソースは、canal-jsondebezium-json (デフォルト)、json 形式をサポートしています。

  • debezium-json を使用する場合、Debezium JSON メッセージには主キー情報が含まれていないため、変換ルールを使用して手動で主キーを追加する必要があります。

    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
  • 単一テーブルのデータが複数のパーティションにまたがっている場合や、パーティションをまたいでテーブルをマージする必要がある場合は、debezium-json.distributed-tables または canal-json.distributed-tablestrue に設定します。

  • Kafka ソースは、複数のスキーマ推論戦略をサポートしています。schema.inference.strategy パラメーターを使用して、優先する戦略を設定します。詳細については、「Kafka」をご参照ください。

リアルタイム Kafka ログの DLF への同期

Kafka クラスターがカスタム JSON 形式でデータを保存している場合、Flink CDC YAML ジョブを設定して DLF に同期できます。システムは、データ型の推論、スキーマ推論、およびスキーマ進化を自動的に処理します。

inventory トピックに、単一のログテーブルのデータが JSON 形式で保存されているとします。このジョブは、そのデータを対応する DLF テーブルに同期します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
  json.infer-schema.flatten-nested-columns.enable: true
  # (オプション) 最初の 100 件の解析エラーをスキップします。エラーが 100 件を超えるとジョブは失敗します。
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# テーブルに主キーを追加します。
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# inventory トピックのすべてのデータを test_database.inventory に書き込みます。
route:
  - source-table: inventory
    sink-table: test_database.inventory

pipeline:
  # (オプション) 処理例外を引き起こしたダーティデータをログに記録します。
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

inventory トピックに、複数のログテーブルのデータが JSON 形式で保存されており、データベース名とテーブル名が databaseName フィールドと tableName フィールドにあるとします。このジョブは、これらのテーブルのデータを対応する DLF ターゲットに同期します。

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
  json.infer-schema.flatten-nested-columns.enable: true
  # databaseName と tableName フィールドの値をデータベース名とテーブル名として使用します。
  json.decode.parser-table-id.fields: databaseName,tableName
  # (オプション) 最初の 100 件の解析エラーをスキップします。エラーが 100 件を超えるとジョブは失敗します。
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# テーブルに主キーを追加します。
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# ods.inventory、ods.customer、ods.user をそれぞれ test_database.inventory、test_database.customer、test_database.user に書き込みます。
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user

pipeline:
  # (オプション) 処理例外を引き起こしたダーティデータをログに記録します。
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

JSON 形式の Kafka ソーステーブルのスキーマ推論と進化戦略の詳細については、「スキーマ解析と変更同期戦略」をご参照ください。

ジョブ設定の詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。