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 コマンドを使用して作成できます。
コンソール (推奨)
-
[Data Management] ページに移動します。
-
Realtime Compute for Apache Flink コンソールにログインし、管理するワークスペースの[操作]列で[コンソール]をクリックします。
-
左側のナビゲーションペインで、[データ管理] をクリックします。
-
-
[カタログの作成] をクリックし、[MySQL] を選択し、[次へ] をクリックします。
-
パラメータを設定します。
重要これらの設定パラメータは、作成後に変更できません。変更を行うには、カタログを削除して再作成する必要があります。
パラメータ
説明
必須
catalogname
MySQL カタログの名前。
はい
hostname
MySQL データベースの IP アドレスまたはホスト名。
説明異なる VPC 間または インターネット 経由で接続するには、ネットワーク接続を確立する必要があります。詳細については、「ネットワーク接続」をご参照ください。
はい
port
MySQL データベースのポート番号。デフォルト: 3306。
いいえ
default-database
デフォルトの MySQL データベースの名前。
はい
username
MySQL データベースのユーザー名。
はい
password
MySQL データベースのパスワード。
シークレットをプレーンテキストで公開しないように、変数の使用を推奨します。この例では、mysqlpw という名前の変数を使用しています。詳細については、「変数の作成」をご参照ください。
はい
-
[OK] をクリックします。
作成されたカタログは、左側の[カタログ] エリアに表示されます。
SQL コマンド
-
[Scripts] ページに移動します。
-
Realtime Compute for Apache Flink コンソールにログインします。管理するワークスペースの [アクション] 列で、[コンソール] をクリックします。
-
左側のナビゲーションペインで、 をクリックします。
-
-
をクリックして [新規スクリプト] をクリックし、[ファイル名] と [保存場所] を入力してから、[保存] をクリックします。 -
次のコードを入力します。
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 以外の値を格納するために使用することは推奨しません。適切な型マッピングを選択してください。詳細については、「型マッピング」をご参照ください。
いいえ
-
-
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 コマンド
-
[Scripts] ページのエディタで、次のコマンドを入力します。
-- Flink でテーブルスキーマを表示します。フィールドのコメントは表示されません。 DESCRIBE `<catalogname>`.`<dbname>`.`<tablename>`; -- カタログを削除します。 DROP CATALOG `<catalogname>`;説明この操作では、カタログのみが削除され、関連するサービス内の基盤となるテーブルは削除されません。カタログのテーブルを使用している実行中のデプロイメントには影響しません。ただし、デプロイメントを再デプロイまたは再起動すると、テーブルが見つからないためエラーが報告されます。慎重に実行してください。
-
コマンドを選択して右クリックし、[実行] を選択します。
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>`