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

Realtime Compute for Apache Flink:Flink CDC データインジェストジョブの開発

最終更新日:Aug 07, 2026

Realtime Compute for Apache Flink では、YAML ファイルを使用して Flink CDC データインジェストジョブを作成し、ソースからシンクへデータを同期できます。このトピックでは、Flink CDC データインジェストジョブを開発する手順について説明します。

背景情報

YAML 設定を使用すると、複雑な ETL パイプラインを簡単に定義でき、これらは自動的に Flink の実行ロジックに変換されます。Flink CDC データインジェストは、Flink CDC に基づく強力なデータ統合ソリューションを提供します。このソリューションは、データベース全体の同期、単一テーブルの同期、シャーディングされたデータベースとテーブルの同期、新しいテーブルの自動検出、スキーマ変更の処理、およびカスタム計算列を効率的にサポートします。また、ETL 処理、WHERE 句フィルタリング、およびカラムプルーニングもサポートしています。この宣言的なアプローチにより、データ統合プロセスが大幅に簡素化され、効率と信頼性が向上します。

Flink CDC の利点

Realtime Compute for Apache Flink では、Flink CDC データインジェストジョブ、SQL ジョブ、または DataStream ジョブを開発してデータ同期を実行できます。以下のセクションでは、他の 2 つの方法と比較した場合の Flink CDC データインジェストジョブの利点について説明します。

Flink CDC と Flink SQL の比較

Flink CDC データインジェストジョブと SQL ジョブは、データ転送に異なるデータ型を使用します。

  • SQL ジョブは、データを RowData として送信します。各 RowData オブジェクトには独自の変更タイプがあり、主なタイプとして挿入 (+I)、更新前 (-U)、更新後 (+U)、削除 (-D) の 4 種類があります。

  • Flink CDC は SchemaChangeEvent を使用して、テーブルの作成、カラムの追加、テーブルの切り捨てなどのスキーマ変更情報を伝達します。また、DataChangeEvent を使用して、挿入、更新、削除などのデータ変更を伝達します。更新メッセージには変更前と変更後の両方の内容が含まれているため、元の変更データをシンクに書き込むことができます。

次の表に、SQL ジョブに対する Flink CDC データインジェストジョブの利点を示します。

Flink CDC インジェスト

Flink SQL

スキーマの自動検出とデータベース全体の同期

手動の CREATE TABLE と INSERT ステートメントが必要です。

複数のスキーマ変更ポリシーをサポート

スキーマ変更をサポートしない

元のチェンジログを保持

元のチェンジログ構造が損なわれる

複数テーブルの読み書きをサポート

単一テーブルの読み書き

CTAS/CDAS 文と比較して、Flink CDC ジョブは、以下のような、より強力な機能を提供します。

  • 新しいデータの書き込みによる同期のトリガーを待つことなく、アップストリームのスキーマ変更を即座に同期。

  • 元のチェンジログの保持、更新メッセージが分割されないことを保証。

  • TRUNCATE TABLE や DROP TABLE など、より多くの種類のスキーマ変更の同期。

  • 柔軟なテーブルマッピングとシンクテーブル名の定義。

  • 柔軟で設定可能なスキーマ進化の動作。

  • WHERE 句によるデータフィルタリング

  • カラムプルーニングのサポート。

Flink CDC と Flink DataStream の比較

次の表に、DataStream ジョブに対する Flink CDC データインジェストジョブの利点を示します。

Flink CDC インジェスト

Flink DataStream

すべてのスキルレベルのユーザーが利用可能です。

Java と分散システムに関する専門知識が必要です。

基盤となる複雑さを隠蔽し、開発を簡素化します。

Flink フレームワークの知識が必要です。

習得しやすい YAML 形式。

依存関係管理のため、Maven のようなツールの知識が必要です。

既存ジョブの再利用性が高いです。

既存コードの再利用が困難です。

制限事項

  • Flink CDC データインジェストジョブを開発するには、Ververica Runtime (VVR) 11.1 以降を使用します。VVR 8.x を使用する必要がある場合は、VVR 8.0.11 を使用してください。

  • 各ジョブは、1 つのソースと 1 つのシンクのみをサポートします。複数のソースから読み取ったり、複数のシンクに書き込んだりするには、複数の Flink CDC ジョブを作成する必要があります。

  • Flink CDC ジョブは、セッションクラスターにはデプロイできません。

  • Flink CDC データインジェストジョブは、自動チューニングをサポートしていません。

Flink CDC データインジェストコネクタ

Flink CDC データインジェストでサポートされているソースおよびシンクコネクタの詳細については、「サポートされているコネクタ」をご参照ください。

Flink CDC データインジェストジョブの作成

テンプレートから作成

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

  2. 対象のワークスペースの Actions 列で、 [Console] をクリックします。

  3. 左側メニューで、Development > Data Ingestion を選択します。

  4. image をクリックし、 [New Draft with Template] をクリックします。

  5. データ同期テンプレートを選択します。

    現在、MySQL から StarRocks、MySQL から Paimon、および MySQL から Hologres へのテンプレートのみが利用可能です。

  6. [name]、 [location]、 [engine version] などのジョブ情報を入力し、 [OK] をクリックします。

  7. Flink CDC ジョブのソースとシンクの情報を設定します。

    パラメータ設定の詳細は、関連するコネクタのドキュメントをご参照ください。

CTAS/CDAS ジョブから作成

重要
  • ジョブに複数の CXAS 文が含まれている場合、Flink は最初の文のみを検出して変換します。

  • Flink SQL と Flink CDC では、組み込み関数のサポートに違いがあるため、生成された transform ルールはそのままでは機能しない場合があります。必要に応じて確認し、調整する必要があります。

  • ソースが MySQL で、元の CTAS/CDAS ジョブがまだ実行中の場合、競合を防ぐために、Flink CDC データ取り込みジョブでソースの server-id を調整する必要があります。

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

  2. 対象のワークスペースの [Actions] 列で、 [Console] をクリックします。

  3. 左側メニューで、Development > Data Ingestion を選択します。

  4. image をクリックし、 [New Draft from CTAS/CDAS] をクリックします。対象の CTAS または CDAS ジョブを選択し、 [OK] をクリックします。

    選択ページでは、有効な CTAS および CDAS ジョブのみが表示されます。通常の ETL ジョブや構文エラーのあるドラフトは表示されません。

  5. [name]、 [location]、 [engine version] などのジョブ情報を入力し、 [OK] をクリックします。

オープンソース Flink CDC から作成

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

  2. 対象のワークスペースの [Actions] 列で、 [Console] をクリックします。

  3. 左側メニューで、Development > Data Ingestion を選択します。

  4. image をクリックし、 [New Draft] を選択します。 [name] と [engine version] を入力し、 [Create] をクリックします。

  5. オープンソースの Flink CDC ジョブの YAML コードをエディターに貼り付けます。

  6. (任意) [Validate] をクリックします。

    このオプションは、構文エラー、ネットワーク接続の問題、および権限の問題をチェックします。

ゼロから作成

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

  2. 対象のワークスペースの [Actions] 列で、 [Console] をクリックします。

  3. 左側メニューで、Development > Data Ingestion を選択します。

  4. image をクリックし、 [New Draft] を選択します。 [name] と [engine version] を入力し、 [Create] をクリックします。

  5. YAML を使用して Flink CDC ジョブを設定します。例:

    # 必須
    source:
      # データソースのタイプ。
      type: <ソースコネクタのタイプに置き換えてください>
      # データソースの設定。設定項目の詳細は、対応するコネクタのドキュメントをご参照ください。
      ...
    # 必須
    sink:
      # シンクのタイプ。
      type: <シンクコネクタのタイプに置き換えてください>
      # シンクの設定。設定項目の詳細は、対応するコネクタのドキュメントをご参照ください。
      ...
    # 任意
    transform:
      # flink_test.customers テーブルの transform ルール。
      - source-table: flink_test.customers
        # プロジェクション設定。同期するカラムを指定して、データ変換を実行します。
        projection: id, username, UPPER(username) as username1, age, (age + 1) as age1, test_col1, __schema_name__ || '.' || __table_name__ identifier_name
        # フィルター条件。id が 10 より大きいデータのみを同期します。
        filter: id > 10
        # transform ルールの説明。
        description: append calculated columns based on source table
    # 任意
    route:
      # ソーステーブルとシンクテーブルのマッピングを指定する route ルール。
      - source-table: flink_test.customers
        sink-table: db.customers_o
        # route ルールの説明。
        description: sync customers table
      - source-table: flink_test.customers_suffix
        sink-table: db.customers_s
        # route ルールの説明。
        description: sync customers_suffix table
    # 任意
    pipeline:
      # ジョブの名前。
      name: MySQL to Hologres Pipeline
    説明

    Flink CDC ジョブでは、キーと値をコロンとスペースで区切ります。形式は Key: Value です。

    次の表に、コードブロックについて説明します。

    必須

    モジュール

    説明

    はい

    source

    データパイプラインの開始点。Flink CDC は、ソースから変更データをキャプチャします。

    説明
    • 現在、サポートされているソースは MySQL のみです。設定の詳細については、「MySQL コネクタ」をご参照ください。

    • 変数を使用して機密情報を管理できます。詳細については、「変数管理」をご参照ください。

    sink

    データパイプラインのエンドポイント。Flink CDC は、キャプチャしたデータ変更をシンクシステムに転送します。

    説明
    • サポートされているシンクについては、「Flink CDC データインジェストコネクタ」をご参照ください。設定の詳細については、特定のコネクタのドキュメントをご参照ください。

    • 変数を使用して機密情報を管理できます。詳細については、「変数管理」をご参照ください。

    いいえ

    pipeline

    (データパイプライン)

    パイプライン名など、データパイプラインジョブ全体の基本設定を定義します。

    transform

    Flink パイプラインを流れるデータを操作するためのデータ変換ルールを指定します。ETL 処理、WHERE 句フィルタリング、列プルーニング、および計算列をサポートします。

    Flink CDC によってキャプチャされた生の変更データを、特定のダウンストリームシステムに合わせて変換するには、transform ブロックを使用します。

    route

    このモジュールが設定されていない場合、ジョブはデフォルトでデータベース全体または対象テーブルの同期を実行します。

    場合によっては、特定のルールに基づいて、キャプチャされた変更データを異なる宛先に送信する必要があります。route モジュールでは、ソースとシンクのマッピング関係を柔軟に指定して、データを異なるターゲットに送信できます。

    各モジュールの構文と設定の詳細は、「Flink CDC データインジェストジョブのリファレンス」をご参照ください。

    次のコードは、MySQL の app_db データベースのすべてのテーブルを Hologres のデータベースに同期する例を示しています。

    source:
      type: mysql
      hostname: <ホスト名>
      port: 3306
      username: ${secret_values.mysqlusername}
      password: ${secret_values.mysqlpassword}
      tables: app_db.\.*
      server-id: 5400-5404
      # (任意) 増分フェーズで新しく作成されたテーブルからデータを同期します。
      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: hologres
      name: Hologres Sink
      endpoint: <エンドポイント>
      dbname: <データベース名>
      username: ${secret_values.holousername}
      password: ${secret_values.holopassword}
    pipeline:
      name: Sync MySQL Database to Hologres
  6. (任意) [Validate] をクリックします。

    このオプションは、構文エラー、ネットワーク接続の問題、および権限の問題をチェックします。

関連トピック