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

Realtime Compute for Apache Flink:Flink CDC ジョブの複雑なユースケース

最終更新日:Jun 22, 2026

このトピックでは、複雑なビジネスシナリオにおける Flink CDC データインジェストジョブのベストプラクティスを紹介します。ソーステーブルのスキーマ変更の処理、メタデータと計算列によるデータの拡張、ソフトデリートの実装、シャード化されたテーブルのマージ、データベース同期、テーブルのフィルタリング、特定のタイムスタンプからのジョブの開始について説明します。

新規テーブルの同期

Flink CDC データインジェストジョブでは、新規テーブルを次の 2 つの方法で同期できます。

  • 空のテーブルのホット同期:履歴データを持たない新規テーブルを動的にキャプチャします。ジョブはそれ以降の変更のみをキャプチャし、再起動は不要です。

  • 履歴データを含むテーブルの同期:既にデータが含まれている新規テーブルに対して全量 + 増分同期を実行します。この場合、ジョブの再起動が必要です。

新規空テーブルのホット同期

増分フェーズ中に新規作成された空のテーブルを再起動なしでリアルタイムに同期するには、scan.binlog.newly-added-table.enabled パラメータを設定します。この方法はジョブの再起動を回避できるため推奨されます。

例えば、MySQL のdlf_testデータベースからすべてのテーブルを同期するデータインジェストジョブを実行しているとします。ソースで products という名前の新規の空テーブルが作成されました。ジョブを再起動せずにこの新規テーブルを同期するには、次のようにジョブ設定で scan.binlog.newly-added-table.enabled: true を設定します。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 増分フェーズ中に新しく作成されたテーブルからデータを同期します。
  scan.binlog.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

この設定により、ジョブはdlf_testデータベースのすべての新規テーブルをシンクに自動的に作成します。

重要

scan.binlog.newly-added-table.enabled パラメータは、scan.startup.modeinitial (デフォルト) に設定されている場合にのみ有効です。

履歴データを含むテーブルの同期

MySQL データベースに既にcustomersテーブルとproductsテーブルが含まれているとします。ただし、最初はcustomersテーブルのみを同期するようにジョブを設定していました。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.customers
  server-id: 8601-8604
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

ジョブがしばらく実行された後、履歴データを含むデータベースのすべてのテーブルを同期する必要がある場合は、ジョブを再起動する必要があります。次の手順に従ってください。

  1. セーブポイントを使用してジョブを停止します。

  2. MySQL ソース設定のtablesパラメータを、すべてのテーブルが対象となるように変更します。次に、scan.binlog.newly-added-table.enabled パラメータが存在する場合は削除し、scan.newly-added-table.enabled を追加します。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true
  1. セーブポイントからジョブを再起動します。

重要

scan.binlog.newly-added-table.enabledscan.newly-added-table.enabled を同時に有効にすることはできません。

特定のテーブルの除外

Flink CDC データインジェストジョブでは、特定のテーブルを同期対象から除外できます。

例えば、dlf_test MySQL データベース内の products_tmp テーブルを除くすべてのテーブルを同期するには、以下の設定を使用します。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  # (オプション) 同期したくないテーブルを除外します。
  tables.exclude: dlf_test.products_tmp
  server-id: 8601-8604
  # (オプション) 増分フェーズ中に新しく作成されたテーブルからデータを同期します。
  scan.binlog.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

この設定により、Flink CDC データインジェストジョブは dlf_test データベースのすべてのテーブルをシンクに自動的に作成しますが、products_tmp テーブルは除外されます。ジョブは、同期対象テーブルのスキーマとデータをリアルタイムで更新し続けます。

説明

tables.excludeパラメータでは、複数のテーブルに一致する正規表現を使用できます。テーブルがtablestables.excludeの両方のパターンに一致する場合、除外ルールが優先され、そのテーブルは同期されません。

メタデータと計算列の追加

メタデータ列の追加

データを書き込む際に、transform モジュールを使用してメタデータ列を追加できます。たとえば、次の設定では、ダウンストリームテーブルにテーブル名、操作時間、操作タイプが追加されます。詳細については、「transform モジュール」をご参照ください。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # 操作時刻をメタデータとして含めます。
  metadata-column.include-list: op_ts
transform:
  - source-table: dlf_test.customers
    projection: __schema_name__ || '.' || __table_name__  as identifier, op_ts, __data_event_type__ as op, *
    # (オプション) プライマリキーを変更します。
    primary-keys: id,identifier
    description: add identifier, op_ts and op
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true
説明

MySQL をソースとして使用する場合、操作時刻をメタデータとしてデスティネーションに送信するには、metadata-column.include-list: op_ts を追加する必要があります。 詳細については、「MySQL」をご参照ください。

ソーステーブルには、すべての変更データイベントタイプが含まれます。シンクテーブルでDELETE操作をINSERTに変換してソフトデリートを実装するには、transformモジュールにconverter-after-transform: SOFT_DELETE設定を追加してください。

計算列の追加

データを書き込む際、transformモジュールを使用して計算列を追加できます。例えば、次の設定では、created_atフィールドを変換してdtフィールドを作成し、それをシンクテーブルのパーティションキーとして使用します。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # 操作時刻をメタデータとして含めます。
  metadata-column.include-list: op_ts
transform:
  - source-table: dlf_test.customers
    projection: DATE_FORMAT(created_at, 'yyyyMMdd') as dt, *
    # (オプション) パーティションキーを設定します。
    partition-keys: dt
    description: add dt
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true
説明

MySQL をソースとして使用する場合、操作時間をメタデータとしてデスティネーションに送信するために、metadata-column.include-list: op_ts を追加する必要があります。 詳細については、「MySQL」をご参照ください。

テーブル名マッピング

routeモジュールを使用して、同期中にテーブル名を変更できます。次の例は、典型的な名前変更シナリオとそれに対応するジョブ設定を示しています。

シャード化されたテーブルのマージ

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
route:
  # dlf_test データベースから product_[0-9]+ パターンに一致するすべてのテーブルを dlf.products テーブルにマージします。
  - source-table: dlf_test.product_[0-9]+
    sink-table: dlf.products
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

データベース同期

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
route:
  # dlf_test データベースのすべてのテーブルを dlf データベースに同期し、テーブル名部分に ods_ プレフィックスを付けて名前を変更します。
  - source-table: dlf_test.\.*
    sink-table: dlf.ods_<>
    replace-symbol: <>
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

総合的なユースケース

次の Flink CDC データインジェストジョブは、このトピックで説明した機能を組み合わせた総合的なユースケースを示しています。このコードを特定のビジネス要件に合わせて調整できます。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  # (オプション) 同期したくないテーブルを除外します。
  tables.exclude: dlf_test.products_tmp
  server-id: 8601-8604
  # (オプション) 新しく追加されたテーブルの全量および増分データを同期します。
  scan.newly-added-table.enabled: true
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
  # 操作時刻をメタデータとして含めます。
  metadata-column.include-list: op_ts
transform:
  - source-table: dlf_test.customers
    projection: __schema_name__ || '.' || __table_name__ as identifier, op_ts, __data_event_type__ as op, DATE_FORMAT(created_at, 'yyyyMMdd') as dt, *
    # (オプション) プライマリキーを変更します。
    primary-keys: id,identifier
    # (オプション) パーティションキーを設定します。
    partition-keys: dt
    # (オプション) DELETE 操作を INSERT 操作に変換してソフトデリートを実現します。
    converter-after-transform: SOFT_DELETE
route:
  # dlf_test データベースのすべてのテーブルを dlf データベースに同期し、テーブル名部分に ods_ プレフィックスを付けて名前を変更します。
  - source-table: dlf_test.\.*
    sink-table: dlf.ods_<>
    replace-symbol: <>
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true

特定のタイムスタンプからの開始

Flink CDC データインジェストジョブをステートレスで起動する際、ソースが特定のバイナリログ (binlog) 位置からデータの読み取りを再開できるよう、開始時刻を指定できます。

運用保守ページでの設定

ジョブの運用保守ページで、ステートレスで起動する際にソーステーブルの開始時刻を指定できます。

この設定は、MySQL および Kafka ソースに対応しています。対応するスイッチを有効にし、特定の開始日時を設定します。

ジョブパラメータの設定

ジョブ定義では、パラメータを設定してソーステーブルの開始時刻を指定できます。

例えば、MySQL ソースの場合、ジョブ設定で scan.startup.mode: timestamp を設定することで、特定のタイムスタンプから開始できます。次に設定例を示します。

source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: username
  password: password
  tables: dlf_test.\.*
  server-id: 8601-8604
  # (オプション) タイムスタンプモードでジョブを開始します。
  scan.startup.mode: timestamp
  # このモードで起動タイムスタンプを指定します。
  scan.startup.timestamp-millis: 1667232000000
  # (オプション) テーブルと列のコメントを同期します。
  include-comments.enabled: true
  # (オプション) 無制限チャンクを優先的にディスパッチし、潜在的な TaskManager の OutOfMemory 問題を防ぎます。
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (オプション) 解析フィルタを有効にして読み取りを高速化します。
  scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  # (オプション) コミットユーザーを指定します。競合を回避するため、ジョブごとに異なるユーザーを使用することを推奨します。
  commit.user: your_job_name
  # (オプション) 削除ベクトルを有効にして読み取りパフォーマンスを向上させます。
  table.properties.deletion-vectors.enabled: true
説明

運用保守ページとジョブパラメータの両方で開始時刻を指定した場合、運用保守ページでの設定が優先されます。