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

Realtime Compute for Apache Flink:リアルタイムデータベースのデータインジェスト

最終更新日:Aug 07, 2026

Realtime Compute for Apache Flink は、フル同期から増分同期への自動切り替え、メタデータ検出、スキーマエボリューション、データベース全体の同期を自動的に処理することで、リアルタイムのデータインジェストを簡素化します。このトピックでは、ApsaraDB RDS for MySQL から Hologres へデータをストリーミングするデータインジェストジョブを迅速に構築する方法を説明します。

背景情報

DMS コンソールでは、これらのデータベースとテーブルは以下のように表示されます。

RDS インスタンスには、複数のデータベース (tpc_ds、tpc_ds_large、user_db1、user_db2、user_db3、__recycle_bin__) が含まれています。user_db3 データベースには、user03、user06、user09 の 3 つのテーブルが含まれています。user03 テーブルには、id (int(11)) と name (varchar(255)) の 2 つの列があります。

以下の手順に従って、これらすべてのテーブルを Hologres に同期し、シャード化されたユーザーテーブルを単一のテーブルにマージするデータインジェストジョブを開発します。

このトピックでは、Flink CDC データインジェストを使用して、データベース全体の同期とシャード化されたテーブルのマージを実行します。これにより、単一のジョブで、フル同期と増分同期、およびリアルタイムのスキーマ変更同期を完了できます。

前提条件

MySQL テストデータと Hologres データベースの準備

  1. tpc_ds.sql、user_db1.sql、user_db2.sql、user_db3.sql をクリックして、テストデータファイルをローカルマシンにダウンロードします。

  2. DMS コンソールで、ApsaraDB RDS for MySQL インスタンスにテストデータを準備します。

    1. DMS を使用して ApsaraDB RDS for MySQL インスタンスにログインします。

    2. SQL コンソール ウィンドウに、次のコマンドを入力し、[実行] をクリックします。

      次のコマンドは、tpc_ds、user_db1、user_db2、user_db3 の 4 つのデータベースを作成します。

      CREATE DATABASE tpc_ds;
      CREATE DATABASE user_db1;
      CREATE DATABASE user_db2;
      CREATE DATABASE user_db3;
    3. トップナビゲーションバーで、[データインポート] をクリックします。

    4. [一括データインポート] タブで、データをインポートするデータベースを選択し、対応する SQL ファイルをアップロードして [送信] をクリックした後、[変更を実行] をクリックします。表示されたダイアログボックスで、[実行を確認] をクリックします。

      この操作を繰り返して、対応するデータファイルを tpc_ds、user_db1、user_db2、user_db3 データベースにインポートします。

      [ファイルエンコーディング] を [自動検出] に、[インポートモード] を [高速モード] または [セーフモード] に、[ファイルタイプ] を [SQL スクリプト]、[CSV]、または [Excel] に設定します。添付ファイルは txt/sql/csv/xlsx/zip 形式に対応しており、サイズは最大 5 GB です。

  3. Hologres コンソールで、マージされたユーザーテーブルデータを格納するために、my_user という名前のデータベースを作成します。

    詳細については、「データベースの作成」をご参照ください。

IP ホワイトリストの設定

Flink が ApsaraDB RDS for MySQL および Hologres インスタンスにアクセスできるようにするには、Flink ワークスペースの CIDR ブロックを両方のインスタンスの IP ホワイトリストに追加します。

  1. Flink ワークスペースの CIDR ブロックを取得します。

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

    2. ワークスペースのリストで、目的の [ワークスペース] を見つけ、[アクション] 列で [その他] > [ワークスペース詳細] を選択します。

    3. [ワークスペース詳細] ダイアログボックスで、 Flink [vSwitch] の [CIDR ブロック]情報を確認します。

  2. Flink の CIDR ブロックを ApsaraDB RDS for MySQL インスタンスの IP ホワイトリストに追加します。

    詳細については、「IP ホワイトリストの設定」をご参照ください。

    [ホワイトリストグループの変更] ダイアログボックスで、[グループホワイトリスト] フィールドに Flink フルマネージドの CIDR ブロックを入力し、[OK] をクリックします。

  3. Flink の CIDR ブロックを Hologres インスタンスの IP ホワイトリストに追加します。

    HoloWeb でデータ接続を設定する際、接続の IP ホワイトリストを設定する前に、[Login Method] を [Password-free login for current user] に設定します。詳細については、「IP ホワイトリスト」をご参照ください。

    HoloWeb で、上部のナビゲーションバーから [セキュリティセンター] を選択し、左側のペインで [IP ホワイトリスト] をクリックし、[IP ホワイトリストの編集] ダイアログボックスで次のパラメーターを設定します。

    • [グループ]: デフォルトを選択します

    • [データベース制限]:すべてを選択

    • [ユーザー制限]:すべて選択

    • [IP アドレス]:Flink CIDR ブロックを CIDR 表記で入力します(例: 172.xx.0/19)。 複数の IP アドレスは改行で区切ります。

    適用するには、[OK] をクリックします。

手順1:データインジェストジョブの開発

  1. Flink 開発コンソールにログインし、新しいジョブを作成します。

    1. [開発] > [データ取り込み] ページで、[新規作成] をクリックします。

    2. [空白のデータ取り込みドラフト] をクリックします。

      Flink は、豊富なコードテンプレートセットを提供しており、それぞれに特定のユースケース、コードサンプル、ガイダンスが含まれています。テンプレートをクリックすると、ビジネスロジックを実装するための製品機能や構文について学習できます。

    3. [次へ] をクリックします。

    4. [新規ドラフト] ダイアログボックスで、パラメーターを設定します。

      パラメーター

      説明

      例

      [名前]

      ジョブの名前。

      説明

      ジョブ名は、現在のプロジェクト内で一意である必要があります。

      flink-test

      [場所]

      ジョブのコードファイルが保存されるフォルダー。

      既存のフォルダーの横にある 新建文件夹 アイコンをクリックして、サブフォルダーを作成することもできます。

      Job Drafts

      [エンジンバージョン]

      ジョブで使用される Flink エンジンのバージョン。エンジンのバージョン番号、バージョンの互換性、および重要なライフサイクル日については、「エンジンバージョン」をご参照ください。

      vvr-11.1-jdk11-flink-1.20

    5. [OK] をクリックします。

  2. 次のジョブコードをジョブエディターにコピーします。

    次のコードは、tpc_ds データベースのすべてのテーブルを Hologres に同期し、シャード化されたユーザーテーブルを Hologres の単一のテーブルにマージします。

    source:
      type: mysql
      name: MySQL Source
      hostname: localhost
      port: 3306
      username: username
      password: password
      tables: tpc_ds.\.*,user_db[0-9]+.user[0-9]+
      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: hologres
      name: Hologres Sink
      endpoint: ****.hologres.aliyuncs.com:80
      dbname: cdcyaml_test
      username: ${secret_values.holo-username}
      password: ${secret_values.holo-password}
      sink.type-normalize-strategy: BROADEN
      
    route:
      # シャード化されたユーザーテーブルを my_user.users テーブルにマージして同期します。
      - source-table: user_db[0-9]+.user[0-9]+
        sink-table: my_user.users
    説明

    MySQL の tpc_ds データベースのテーブルは、同期先で同じ名前のテーブルに直接マッピングされるため、route セクションに追加のマッピング設定は必要ありません。ods_tps_ds のような別の名前のデータベースにテーブルを同期するには、route モジュールを次のように設定します。

    route:
      # シャード化されたユーザーテーブルを my_user.users テーブルにマージして同期します。
      - source-table: user_db[0-9]+.user[0-9]+
        sink-table: my_user.users
      # tpc_ds 配下のすべてのテーブルのデータベース名を変更し、ods_tps_ds に同期します。
      - source-table: tpc_ds.\.*
        sink-table: ods_tps_ds.<>
        replace-symbol: <>

手順2:ジョブの開始

  1. [開発] > [データ取り込み] ページで、[デプロイ] をクリックします。表示されるダイアログボックスで、[確認] をクリックします。

    ダイアログボックスで、[備考] を入力し、[ジョブタグ] を設定して、[デプロイターゲット] ドロップダウンリストからキューを選択できます (デフォルト: default-queue)。 注:デプロイは、次回のジョブ起動時に有効になります。

  2. [O&M] > [デプロイメント] ページで、対象のジョブの [アクション] 列にある [開始] をクリックします。必要に応じてパラメーターを設定します。詳細については、「ジョブを開始する」をご参照ください。

  3. [開始] をクリックします。

    ジョブが開始されると、「ジョブオペレーション」ページでそのランタイム情報とステータスを表示できます。ジョブのステータスには、[実行中]、[失敗]、[停止] があります。上部の [ストリームジョブ] ドロップダウンフィルターを使用して、ジョブリストをタイプ別にフィルターできます。

手順3:フル同期結果の検証

  1. Hologres 管理コンソールにログインします。

  2. [メタデータ管理] タブで、Hologres インスタンスの tpc_ds データベースに 24 個のテーブルとそのデータが存在することを確認します。

    インスタンス > [tpc_ds] > [public] > [テーブル] に移動すると、call_center、catalog_sales、customer、store_sales などのテーブルが表示されます。 [store_sales] テーブルを選択し、[データプレビュー] タブをクリックすると、ss_sold_date、ss_sold_time、ss_item_sk、ss_customer などの列が表示されます。

  3. [メタデータ管理] タブで、my_user データベースの users テーブルのスキーマを確認します。

    同期されたテーブルのスキーマとデータは以下のとおりです。

    • テーブルスキーマ

      同期後、users テーブルには次の列が含まれます。

      • _db_name (text, プライマリキー)

      • _table_name (text, プライマリキー)

      • id (int4, プライマリキー)

      • name (varchar 255, null 許容)

      _db_name、_table_name、id の列が複合プライマリキーを構成します。

      users テーブルのスキーマには、ソースの MySQL テーブルにはない 2 つの追加列、_db_name と _table_name が含まれています。これらの列は各行のソースデータベースとソーステーブルを示し、複合プライマリキーの一部として、シャード化されたテーブルがマージされた後のデータの一意性を保証します。

    • テーブルデータ

      users テーブル情報ページの右上隅で、[テーブルのクエリ] をクリックします。次のコマンドを入力し、[実行] をクリックします。

      select * from users order by _db_name,_table_name,id;

      クエリ結果から、ユーザーデータが 3 つのデータベース (user_db1 (user01、user04、user07)、user_db2 (user02、user05、user08)、user_db3 (user03、user06、user09)) に完全に同期され、_db_name、_table_name、id、name の列を持つ合計 9 件のレコードがあることがわかります。

手順4:増分同期結果の検証

フル同期が完了すると、ジョブは手動介入なしで自動的に増分同期フェーズに切り替わります。[Monitoring and Alerts] タブで currentEmitEventTimeLag の値を確認することで、データ同期フェーズを判断できます。

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

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

  3. [O&M] > [デプロイメント] ページで、対象のジョブの名前をクリックします。

  4. [アラーム] (または [メトリクス]) タブをクリックします。

  5. currentEmitEventTimeLag チャートを調べて、データ同期フェーズを判断します。

    数据曲线

    • 値が 0 の場合は、フル同期フェーズを示します。

    • 値が 0 より大きい場合は、増分同期フェーズを示します。

  6. リアルタイムのデータとスキーマ変更の同期を検証します。

    MySQL CDC ソースは、増分フェーズ中のリアルタイムのデータとスキーマの同期をサポートしています。これを検証するには、ジョブがこのフェーズに入った後、MySQL でシャード化されたユーザーテーブルのスキーマとデータを変更します。

    1. DMS を使用して ApsaraDB RDS for MySQL インスタンスにログインします。

    2. user_db2 データベースで、次のコマンドを実行して user02 テーブルのスキーマを変更し、データを挿入および更新します。

      USE `user_db2`;
      ALTER TABLE `user02` ADD COLUMN `age` INT;   -- age 列を追加します。
      INSERT INTO `user02` (id, name, age) VALUES (27, 'Tony', 30); -- age データを含む行を挿入します。
      UPDATE `user05` SET name='JARK' WHERE id=15;  -- 別のテーブルを更新し、名前を大文字に変更します。
    3. Hologres コンソールで、users テーブルのスキーマとデータの変更を確認します。

      users テーブル情報ページの右上隅で、[テーブルをクエリ] をクリックして次のコマンドを入力し、[実行] をクリックします。

      select * from users order by _db_name,_table_name,id;

      シャード化されたテーブルのスキーマが異なっていても、user02 のスキーマ変更とデータ変更はリアルタイムで伝播されます。Hologres の users テーブルには、新しい age 列、挿入された Tony のレコード、および更新された JARK のレコードが表示されます。

      クエリ結果には、user_db1、user_db2、user_db3 にまたがる users テーブルから 10 行が表示されます。増分同期の変更が反映されています。user_db2 のユーザー Tony (id 27) は age が 30 に更新され、別のレコードは name が JARK に変更されています。残りの行の age の値は \N であり、増分同期が機能していることを確認できます。

(任意) 手順5:ジョブリソースの設定

パフォーマンスを向上させるために、データ量に基づいて、並列度、TaskManager のメモリ、CU などのジョブリソースを調整できます。

  1. [O&M] > [デプロイメント] ページで、対象のジョブの名前をクリックします。

  2. [設定] タブで、[リソース] セクションの右上隅にある [編集] をクリックします。

  3. TaskManager のメモリや並列度などのリソースパラメーターを手動で設定します。

  4. [リソース] セクションの右側で、[保存] をクリックします。

  5. ジョブを再起動します。

    リソース設定の変更は、ジョブを再起動した後にのみ有効になります。

関連ドキュメント