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

Realtime Compute for Apache Flink:MySQL カタログの管理

最終更新日:Sep 19, 2026

MySQL カタログを使用すると、Realtime Compute for Apache Flink コンソールから MySQL インスタンス内のテーブルに直接アクセスし、Flink SQL デプロイメントで使用できます。このトピックでは、MySQL カタログの作成方法と使用方法について説明します。

背景情報

MySQL カタログには、次の機能があります。

  • DDL ステートメントを使用した手動登録なしで、MySQL インスタンス内のテーブルに直接アクセスできるため、開発効率と精度が向上します。

  • MySQL カタログのテーブルを、Flink SQL デプロイメントで CDC ソーステーブル、シンクテーブル、またはディメンションテーブルとして使用できます。

  • ApsaraDB RDS for MySQL、PolarDB for MySQL、およびセルフマネージド MySQL データベースをサポートしています。

  • シャードテーブルの論理テーブルへの直接アクセスをサポートしています。

  • Flink CDC データインジェストジョブと連携して、MySQL データソースのデータベース全体の同期、シャードテーブルのマージ同期、およびスキーマ変更の同期を実行します。

制限事項

  • Realtime Compute for Apache Flink と MySQL インスタンスは同じ VPC 内にある必要があります。異なる VPC 間またはインターネット経由で接続するには、ネットワーク接続を確立する必要があります。詳細については、「ネットワーク接続」をご参照ください。

  • カタログの設定は、作成後に変更できません。設定を変更するには、カタログを削除して再作成する必要があります。

  • Flink でデータベースやテーブルを作成することはできません。

  • ソースとして使用する場合、これらのテーブルはストリーム読み取りのみをサポートし、バッチ読み取りはサポートしていません。

    説明

    MySQL カタログのテーブルを CDC ソーステーブルとして使用する前に、ApsaraDB RDS for MySQL、PolarDB for MySQL、またはセルフマネージド MySQL データベースでバイナリロギング (Binlog) を有効にする必要があります。詳細については、「MySQL データベースの設定」をご参照ください。

  • カタログは、DDL ステートメントで PolarDB 固有の構文を使用しているテーブルを識別できません。

    例: PARTITION BY KEY(`idempotent_id`) PARTITIONS 16, UNIQUE KEY `uk_order_id` (`order_id`)

  • Ververica Runtime (VVR) 8.0.7 以降では、ビューを Flink テーブルとして使用できません。

  • MySQL バージョン 5.7 および 8.0.x のみがサポートされています。

MySQL カタログの作成

MySQL カタログは、コンソールまたは SQL コマンドを使用して作成できます。

コンソール (推奨)

  1. [Data Management] ページに移動します。

    1. Realtime Compute for Apache Flink コンソールにログインし、管理するワークスペースの[操作]列で[コンソール]をクリックします。

    2. 左側のナビゲーションペインで、[データ管理] をクリックします。

  2. [カタログの作成] をクリックし、[MySQL] を選択し、[次へ] をクリックします。

  3. パラメータを設定します。

    重要

    これらの設定パラメータは、作成後に変更できません。変更を行うには、カタログを削除して再作成する必要があります。

    パラメータ

    説明

    必須

    catalogname

    MySQL カタログの名前。

    はい

    hostname

    MySQL データベースの IP アドレスまたはホスト名。

    説明

    異なる VPC 間または インターネット 経由で接続するには、ネットワーク接続を確立する必要があります。詳細については、「ネットワーク接続」をご参照ください。

    はい

    port

    MySQL データベースのポート番号。デフォルト: 3306。

    いいえ

    default-database

    デフォルトの MySQL データベースの名前。

    はい

    username

    MySQL データベースのユーザー名。

    はい

    password

    MySQL データベースのパスワード。

    シークレットをプレーンテキストで公開しないように、変数の使用を推奨します。この例では、mysqlpw という名前の変数を使用しています。詳細については、「変数の作成」をご参照ください。

    はい

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

    作成されたカタログは、左側の[カタログ] エリアに表示されます。

SQL コマンド

  1. [Scripts] ページに移動します。

    1. Realtime Compute for Apache Flink コンソールにログインします。管理するワークスペースの [アクション] 列で、[コンソール] をクリックします。

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

  2. image をクリックして [新規スクリプト] をクリックし、[ファイル名] と [保存場所] を入力してから、[保存] をクリックします。

  3. 次のコードを入力します。

    CREATE CATALOG YourCatalogName WITH(
      'type' = 'mysql',
      'hostname' = 'rm-bp1gcn0q0j0******.mysql.rds.aliyuncs.com',
      'port' = '3306',
      'username' = 'usertest',
      'password' = '${secret_values.mysqlpw}',
      'default-database' = 'flinktest',
      'catalog.table.metadata-columns'='table_name'
    );

    パラメータ

    説明

    必須

    YourCatalogName

    MySQL カタログの名前。

    はい

    type

    カタログのタイプ。mysql に設定します。

    はい

    hostname

    MySQL データベースの IP アドレスまたはホスト名。

    説明

    異なる VPC 間または インターネット 経由で接続するには、ネットワーク接続を確立する必要があります。詳細については、「ネットワーク接続」をご参照ください。

    はい

    port

    MySQL データベースのポート番号。デフォルト: 3306。

    いいえ

    default-database

    デフォルトの MySQL データベースの名前。

    はい

    username

    MySQL データベースのユーザー名。

    はい

    password

    MySQL データベースのパスワード。

    シークレットをプレーンテキストで公開しないように、変数の使用を推奨します。この例では、mysqlpw という名前の変数を使用しています。詳細については、「変数の作成」をご参照ください。

    はい

    property-version

    カタログのプロパティスキーマのバージョン。0 または 1 (推奨) に設定します。

    バージョンによって、サポートされるプロパティとデフォルト値が異なる場合があります。詳細については、プロパティの説明をご参照ください。

    説明
    • VVR 8.0.6 以降でのみサポートされています。

    • VVR 11.1 以降では、デフォルト値は 1 です。それ以前のバージョンでは、デフォルト値は 0 です。

    いいえ

    catalog.table.metadata-columns

    テーブルをクエリする際に、テーブルスキーマに追加する MySQL CDC ソーステーブルのメタデータ列を指定します。デフォルトでは、メタデータ列は追加されません。

    複数のメタデータ列はセミコロン (;) で区切ります。例: op_ts;table_name;database_name

    説明
    • VVR 6.0.5 以降でのみサポートされています。

    • このプロパティを設定すると、指定されたメタデータ列がスキーマに追加されます。これらの列は MySQL CDC ソーステーブル固有であるため、このカタログのテーブルはソーステーブルとしてのみ使用でき、シンクテーブルやディメンションテーブルとしては使用できません。

    いいえ

    catalog.table.treat-tinyint1-as-boolean

    テーブルスキーマを取得する際に、MySQL の TinyInt(1) と Boolean を Flink の Boolean にマッピングするかどうかを指定します。有効な値:

    • true:Boolean にマッピングします。

    • false:TINYINT にマッピングします。

    デフォルト値:

    • property-version が 0 の場合、デフォルト値は true です。

    • property-version が 1 の場合、デフォルト値は false です。

    説明
    • VVR 8.0.4 以降でのみサポートされています。

    • MySQL で TinyInt(1) を 0 と 1 以外の値を格納するために使用することは推奨しません。適切な型マッピングを選択してください。詳細については、「型マッピング」をご参照ください。

    いいえ

  4. CREATE CATALOG ステートメントを選択し、左側の行番号の横にある[実行]をクリックします。

    The following statement has been executed successfully! というメッセージが表示されたら、カタログが正常に作成されています。

    エディター内の SQL ステートメント CREATE CATALOG myCatalog には、type=mysql、hostname=rm-bp1gcn0q0j0**.mysql.rds.aliyuncs.com、port=3306、username=usertest、password=${secret_values.mysqlpw} (変数への参照)、default-database=flinktest、および catalog.table.metadata-columns=table_name などの設定パラメーターがあります。

MySQL カタログの表示と削除

コンソール (推奨)

[データ管理] ページでは、[カタログリスト] で、作成されたカタログの [名前] と [タイプ] を表示できます。

  • 表示: カタログの[操作]列で、[表示]をクリックすると、カタログ内のデータベースとテーブルが表示されます。

    テーブルスキーマの詳細には、フィールドのコメントは表示されません。

  • 削除: カタログの [アクション] 列で、[削除] をクリックします。

    この操作では、カタログのみが削除され、関連するサービス内の基盤となるテーブルは削除されません。カタログのテーブルを使用している実行中のデプロイメントには影響しません。ただし、これらのデプロイメントを再デプロイまたは再起動すると、テーブルが見つからないためエラーが報告されます。慎重に実行してください。

SQL コマンド

  1. [Scripts] ページのエディタで、次のコマンドを入力します。

    -- Flink でテーブルスキーマを表示します。フィールドのコメントは表示されません。
    DESCRIBE `<catalogname>`.`<dbname>`.`<tablename>`;
    -- カタログを削除します。
    DROP CATALOG `<catalogname>`;
    説明

    この操作では、カタログのみが削除され、関連するサービス内の基盤となるテーブルは削除されません。カタログのテーブルを使用している実行中のデプロイメントには影響しません。ただし、デプロイメントを再デプロイまたは再起動すると、テーブルが見つからないためエラーが報告されます。慎重に実行してください。

  2. コマンドを選択して右クリックし、[実行] を選択します。

    DESCRIBE ステートメントを実行すると、orderkey、custkey、order_status、total_price などのフィールドとそのデータ型およびプロパティを含む、テーブルのスキーマが返されます。

MySQL カタログの使用

MySQL ソーステーブルからの読み取り

INSERT INTO `<othersinktable>`
SELECT ...
FROM `<mysqlcatalog>`.`<dbname>`.`<tablename>` /*+ OPTIONS('server-id' = '6000-6008') */;

MySQL カタログのテーブルを CDC ソーステーブルとして使用する場合は、SQL ヒントを使用して、各 Flink SQL デプロイメントに異なる server-id を指定することを推奨します。ソーステーブルがより高い並列度を必要とする場合は、server-id を並列度以上の大きさの範囲として設定する必要があります。

シャードテーブルの論理テーブルからの読み取り

MySQL カタログは、データベース名とテーブル名に正規表現を使用して、シャードテーブルのデータを単一の論理テーブルとして読み取ることをサポートしています。

たとえば、シャードデータベースには、db01 から db10 などのデータベースに分散された user01、user02、user99 などの複数のテーブルが含まれています。すべてのテーブルに互換性のあるスキーマがある場合、データベース名とテーブル名に正規表現を使用して、すべての user シャードテーブルにアクセスできます。

SELECT ... FROM `db.*`.`user.*` /*+ OPTIONS('server-id'='6000-6018') */;

シャードテーブルの論理テーブルは、_db_name (STRING) と _table_name (STRING) という 2 つの追加のシステムフィールドを返します。これらのフィールドは、元のプライマリキーと一緒に、一意性を保証する新しい複合プライマリキーを形成します。たとえば、user01 から user99 テーブルのプライマリキーが id の場合、論理 user テーブルの複合プライマリキーは (_db_name, _table_name, id) になります。

MySQL カタログは、正規表現を使用して同期する複数のテーブルを照合することをサポートしており、シャードテーブルのマージ同期を可能にします。例については、「シャードテーブルのマージと同期」をご参照ください。

Flink CDC データインジェストを使用した MySQL データとスキーマ変更のリアルタイム同期

Flink CDC データインジェストは、単一テーブル同期、スキーマ変更の同期、シャードテーブルのマージ同期、およびカスタム計算列を使用した同期をサポートしています。また、スキーマ変更を含む、データベースレベルのスキーマとデータのリアルタイム同期もサポートしています。例と詳細については、「MySQL YAML コネクタ」をご参照ください。

# 単一テーブル同期:テーブルレベルのスキーマ変更とデータ変更をリアルタイムで同期します。
source:
  type: mysql
  using.built-in-catalog: mysql-catalog
  tables: "<db-name>.<table-name>"
  server-id: "6000-6018"
  # (オプション) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
  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: xxx
  xxx

route:
  - source-table: "<db-name>.<table-name>"
    sink-table: "<target-db-name>.<target-table-name>"
# データベース全体の同期:データベースレベルのスキーマ変更とデータ変更をリアルタイムで同期します。
source:
  type: mysql
  using.built-in-catalog: mysql-catalog
  tables: '<db-name>.\.*'
  server-id: "6000-6018"
  # (オプション) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
  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: xxx
  xxx

route:
  - source-table: '<db-name>.\.*'
    sink-table: "<target-db-name>.<>"
    replace-symbol: "<>"

たとえば、MySQL データを Hologres に同期する方法については、「Hologres カタログの使用」をご参照ください。

source:
  type: mysql
  using.built-in-catalog: mysql-catalog
  tables: dbmysql.mysqltable
  server-id: "8001-8004"
  # (オプション) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
  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: hologres
  using.built-in-catalog: hologres-catalog
  jdbcWriteBatchSize: 1024 # オプション。シンクテーブルのパラメータを指定します。

route:
  - source-table: dbmysql.mysqltable
    sink-table: public.holotable

MySQL ディメンションテーブルからの読み取り

INSERT INTO `<othersinktable>`
SELECT ...
FROM `<othersourcetable>` AS e
JOIN `<mysqlcatalog>`.`<dbname>`.`<tablename>` FOR SYSTEM_TIME AS OF e.proctime AS w
ON e.id = w.id;

MySQL テーブルへの書き込み

INSERT INTO `<mysqlcatalog>`.`<dbname>`.`<tablename>`
SELECT ...
FROM `<othersourcetable>`