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

Realtime Compute for Apache Flink:MySQL

最終更新日:Jul 11, 2026

このトピックでは、MySQL コネクタの使用方法を説明します。

背景情報

MySQL コネクタは、ApsaraDB RDS for MySQL、PolarDB for MySQL、OceanBase (MySQL mode)、セルフマネージド MySQL など、MySQL プロトコルと互換性のあるすべてのデータベースをサポートしています。

重要

MySQL コネクタを使用して OceanBase からデータを読み取る場合は、バイナリロギング (binlog) が有効化され、正しく設定されていることを確認してください。詳細については、「Binlog-related operations」をご参照ください。この機能はパブリックプレビューのため、慎重に使用してください。

MySQL コネクタは、次の項目をサポートしています。

カテゴリ

詳細

サポートされるタイプ

ソーステーブル、ディメンションテーブル、シンクテーブル、およびデータインジェストのデータソース

ランタイムモード

ストリーミングモードのみをサポートしています。

データ形式

該当なし

固有の監視メトリクス

監視メトリクス

  • ソーステーブル

    • currentFetchEventTimeLag:データが生成されてから Source オペレーターによってプルされるまでの間隔。

      このメトリクスはバイナリロギングフェーズでのみ有効です。スナップショットフェーズでは、値は常に 0 です。

    • currentEmitEventTimeLag:データが生成されてから Source オペレーターを離れるまでの間隔。

      このメトリクスはバイナリロギングフェーズでのみ有効です。スナップショットフェーズでは、値は常に 0 です。

    • sourceIdleTime:ソーステーブルで新しいデータが生成されない期間。

  • ディメンションテーブルとシンクテーブル:なし。

説明

メトリクスの詳細については、「Metric description」をご参照ください。

API タイプ

DataStream、SQL、およびデータインジェスト YAML

シンクテーブル内のデータの更新または削除のサポート

はい

機能

MySQL 変更データキャプチャ (CDC) ソーステーブルは、MySQL ストリーミングソーステーブルとも呼ばれ、まずデータベースから完全な履歴データを読み取ります。次に、シームレスにバイナリログの読み取りに切り替えます。このプロセスにより、データの欠落や重複がないことが保証されます。障害が発生した場合でも、データは正確に 1 回のセマンティクスで処理されます。MySQL CDC ソーステーブルは、全データの同時読み取りをサポートしています。増分スナップショットアルゴリズムを使用して、ロックフリーの読み取りと再開可能なデータ転送を実装します。詳細については、「MySQL CDC ソーステーブルについて」をご参照ください。

  • 2つの個別プロセスを維持する必要をなくす、全データと増分データの両方を読み取り可能な統合バッチ/ストリーム処理

  • パフォーマンスの水平スケーリングを可能にする、全データの同時読み取り

  • 全データ読み取りから増分データ読み取りへのシームレスな切り替えと、コンピューティングリソースを節約するための自動スケールイン

  • 安定性を向上させる、全データ読み取りフェーズ中の再開可能なデータ転送

  • オンラインサービスに影響を与えない、全データのロックフリーの読み取り

  • ApsaraDB RDS for MySQL のバックアップログ読み取りのサポート

  • 読み取りレイテンシーを低減する、バイナリログファイルの並列解析

前提条件

MySQL CDC ソーステーブルを使用する前に、「Configure MySQL」に記載されている前提条件の操作を完了する必要があります。

ApsaraDB RDS for MySQL

  • ネットワークプローブを実行し、Realtime Compute for Apache Flink へのネットワーク接続を確保してください。

  • MySQL バージョン:5.6、5.7、8.0.x、または 8.4。

  • バイナリロギングが有効になっている必要があります。デフォルトで有効になっています。

  • バイナリログフォーマットは ROW にする必要があります。これはデフォルトのフォーマットです。

  • binlog_row_image パラメータを FULL に設定する必要があります。これはデフォルトの設定です。

  • バイナリログトランザクション圧縮を無効にする必要があります。この機能は MySQL 8.0.20 で導入され、デフォルトで無効になっています。

  • MySQL ユーザーを作成し、SELECT、SHOW DATABASES、REPLICATION SLAVE、および REPLICATION CLIENT の権限を付与する必要があります。

  • MySQL データベースとテーブルを作成してください。詳細については、「Create a database and an account for an ApsaraDB RDS for MySQL instance」をご参照ください。権限不足による操作の失敗を防ぐため、特権アカウントを使用して MySQL データベースを作成してください。

  • IP アドレスホワイトリストを設定してください。詳細については、「Configure an IP address whitelist for an ApsaraDB RDS for MySQL instance」をご参照ください。

PolarDB for MySQL

  • ネットワークプローブを実行し、Realtime Compute for Apache Flink へのネットワーク接続を確保してください。

  • MySQL バージョン:5.6、5.7、8.0.x、または 8.4。

  • バイナリロギングが有効になっている必要があります。デフォルトでは無効になっています。

  • バイナリログフォーマットは ROW にする必要があります。これはデフォルトのフォーマットです。

  • binlog_row_image パラメータを FULL に設定する必要があります。これはデフォルトの設定です。

  • バイナリログトランザクション圧縮を無効にする必要があります。この機能は MySQL 8.0.20 で導入され、デフォルトで無効になっています。

  • MySQL ユーザーを作成し、SELECT、SHOW DATABASES、REPLICATION SLAVE、および REPLICATION CLIENT の権限を付与する必要があります。

  • MySQL データベースとテーブルを作成してください。詳細については、「Create a database and an account for a PolarDB for MySQL cluster」をご参照ください。権限不足による操作の失敗を防ぐため、特権アカウントを使用して MySQL データベースを作成してください。

  • IP アドレスホワイトリストを設定してください。詳細については、「Configure an IP address whitelist for a PolarDB for MySQL cluster」をご参照ください。

セルフマネージド MySQL

  • ネットワークプローブを実行し、Realtime Compute for Apache Flink へのネットワーク接続を確保してください。

  • MySQL バージョン:5.6、5.7、8.0.x、または 8.4。

  • バイナリロギングが有効になっている必要があります。デフォルトでは無効になっています。

  • バイナリログフォーマットは ROW にする必要があります。デフォルトのフォーマットは STATEMENT です。

  • binlog_row_image パラメータを FULL に設定する必要があります。これはデフォルトの設定です。

  • バイナリログトランザクション圧縮を無効にする必要があります。この機能は MySQL 8.0.20 で導入され、デフォルトで無効になっています。

  • MySQL ユーザーを作成し、SELECT、SHOW DATABASES、REPLICATION SLAVE、および REPLICATION CLIENT の権限を付与する必要があります。

  • MySQL データベースとテーブルを作成してください。詳細については、「Create a database and an account for a self-managed MySQL instance」をご参照ください。権限不足による操作の失敗を防ぐため、特権アカウントを使用して MySQL データベースを作成してください。

  • IP アドレスホワイトリストを設定してください。詳細については、「Configure an IP address whitelist for a self-managed MySQL instance」をご参照ください。

使用制限

一般的な制限

  • MySQL CDC ソーステーブルは、ウォーターマークの定義をサポートしていません。

  • CTAS (Create Table As Select) および CDAS (Create Database As Select) ジョブでは、MySQL CDC ソーステーブルは一部のスキーマ変更を同期できます。サポートされる変更タイプの詳細については、「Schema evolution synchronization policies」をご参照ください。

  • MySQL CDC コネクタは、バイナリログトランザクション圧縮機能をサポートしていません。そのため、MySQL CDC コネクタを使用して増分データを取得する場合は、バイナリログトランザクション圧縮が無効になっていることを確認してください。無効になっていない場合、コネクタが増分データの取得に失敗する可能性があります。

ApsaraDB RDS for MySQL の制限

  • ApsaraDB RDS for MySQL の場合、セカンダリデータベースまたはリードレプリカからデータを読み取らないでください。これは、セカンダリデータベースとリードレプリカのデフォルトのバイナリログ保持期間が短いためです。バイナリログの有効期限が切れてクリアされると、ジョブがバイナリログデータの取得に失敗し、エラーが発生する可能性があります。

  • ApsaraDB RDS for MySQL は、デフォルトで並列プライマリ/セカンダリ同期を有効にしていますが、プライマリインスタンスとセカンダリインスタンス間で一貫したトランザクション順序を保証していません。これにより、プライマリ/セカンダリ スイッチオーバーおよびチェックポイントリカバリ中にデータが欠落する可能性があります。この問題を回避するには、ApsaraDB RDS for MySQL の slave_preserve_commit_order オプションを手動で有効にすることができます。

PolarDB for MySQL の制限

MySQL CDC ソーステーブルは、PolarDB for MySQL V1.0.19 以前のマルチマスタークラスターアーキテクチャのクラスターからのデータ読み取りをサポートしていません。詳細については、「What is a Multi-master Cluster?」をご参照ください。これらのクラスターによって生成されるバイナリログには、重複したテーブル ID が含まれている可能性があります。これにより、CDC ソーステーブルでスキーママッピングエラーが発生し、バイナリログデータの解析時にエラーが発生する可能性があります。

オープンソース MySQL の制限

デフォルトでは、MySQL はプライマリ/セカンダリ間のバイナリログレプリケーション中にトランザクション順序を維持します。MySQL レプリカでパラレルレプリケーションが有効になっている (slave_parallel_workers > 1) ものの、slave_preserve_commit_order = ON が有効になっていない場合、そのトランザクションコミット順序がプライマリデータベースと一致しない可能性があります。Flink CDC がチェックポイントから回復する際、順序の乱れによりデータが欠落する可能性があります。MySQL レプリカで slave_preserve_commit_order = ON を設定できます。または、 slave_parallel_workers = 1 を設定することもできますが、これによりレプリケーションのパフォーマンスが低下します。

注意事項

  • ソーステーブル

    • 各 MySQL CDC データソースには、一意のサーバー ID が必要です。

      サーバー ID の目的

      各 MySQL CDC データソースには、一意のサーバー ID が必要です。複数の MySQL CDC データソースが同じサーバー ID を共有し、再利用できない場合、バイナリログのオフセットが乱れる可能性があります。これにより、データが複数回読み取られたり、欠落したりする可能性があります。

      異なるシナリオにおけるサーバー ID の設定

      サーバー ID は、データ定義言語 (DDL) ステートメントで指定できます。ただし、DDL パラメータではなく、動的ヒントを使用してサーバー ID を設定することを推奨します。

      • 並列度 = 1 またはインクリメンタルスナップショットが無効な場合

        ## インクリメンタルスナップショットフレームワークが無効になっている場合、または並列度が 1 の場合、特定のサーバー ID を指定できます。
        SELECT * FROM source_table /*+ OPTIONS('server-id'='123456') */ ;
      • 並列度 > 1 かつインクリメンタルスナップショットが有効な場合

        ## サーバー ID の範囲を指定する必要があります。範囲内で使用可能なサーバー ID の数は、並列度以上にする必要があります。並列度が 3 であると仮定します。
        SELECT * FROM source_table /*+ OPTIONS('server-id'='123456-123458') */ ;
      • データ同期のための CTAS

        CTAS を使用してデータを同期する場合、CDC データソースの設定が同じであれば、データソースは自動的に再利用されます。この場合、複数の CDC データソースに同じサーバー ID を設定できます。詳細については、「Example 4: Multiple CTAS statements」をご参照ください。

      • 複数の非 CTAS ソーステーブルが再利用できない場合

        ジョブに複数の MySQL CDC ソーステーブルが含まれており、同期に CTAS ステートメントを使用していない場合、データソースは再利用できません。各 CDC ソーステーブルに異なるサーバー ID を指定する必要があります。同様に、インクリメンタルスナップショットフレームワークが有効になっており、並列度が 1 より大きい場合は、サーバー ID の範囲を指定する必要があります。

        select * from 
          source_table1 /*+ OPTIONS('server-id'='123456-123457') */
        left join 
          source_table2 /*+ OPTIONS('server-id'='123458-123459') */
        on source_table1.id=source_table2.id;
    • 全データ読み取りフェーズ中は、セーブポイントを保存したり、ソーステーブルにテーブルを追加または削除したりした後、セーブポイントからジョブを再起動することはできません。これらの操作を実行すると、ジョブはデータの読み取りに失敗します。

  • シンクテーブル

    • 自動インクリメントプライマリキー:DDL で自動インクリメントプライマリキーを宣言しないでください。MySQL はデータ書き込み時に自動的に値を設定します。

    • 少なくとも 1 つの主キー以外のフィールドを宣言する必要があります。宣言しない場合、エラーが発生します。

    • DDL の NOT ENFORCED 制約は、Flink がプライマリキーの検証を強制しないことを示します。プライマリキーの正確性と整合性を保証するのは、ユーザーの責任です。詳細については、「Validity Check」をご参照ください。

  • ディメンションテーブル

    インデックスを使用してクエリを高速化する場合、JOIN 句のフィールドの順序は、インデックスで定義された順序と一致する必要があります。これは最左接頭辞ルールに基づいています。たとえば、インデックスが (a, b, c) の場合、JOIN 条件は ON t.a = x AND t.b = y となります。

    Flink によって生成された SQL は、オプティマイザによって書き換えられる可能性があります。これにより、実際のデータベースクエリ中にインデックスがヒットしない場合があります。インデックスが使用されているかどうかを確認するには、実行計画 (EXPLAIN) または MySQL のスロークエリログで、実際に実行される SELECT ステートメントを確認してください。

SQL

SQL ジョブでは、MySQL コネクタをソーステーブル、ディメンションテーブル、またはシンクテーブルとして使用できます。

構文

CREATE TEMPORARY TABLE mysqlcdc_source (
   order_id INT,
   order_date TIMESTAMP(0),
   customer_name STRING,
   price DECIMAL(10, 5),
   product_id INT,
   order_status BOOLEAN,
   PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
  'connector' = 'mysql',
  'hostname' = '<yourHostname>',
  'port' = '3306',
  'username' = '<yourUsername>',
  'password' = '<yourPassword>',
  'database-name' = '<yourDatabaseName>',
  'table-name' = '<yourTableName>'
);
説明
  • シンクテーブルに書き込む際、コネクタは受信した各レコードに対して、次のように SQL ステートメントを構築して実行します:

    • プライマリキーのないシンクテーブルの場合、 INSERT INTO table_name (column1, column2, ...) VALUES (value1, value2, ...); ステートメントが実行されます。

    • プライマリキーのあるシンクテーブルの場合、 INSERT INTO table_name (column1, column2, ...) VALUES (value1, value2, ...) ON DUPLICATE KEY UPDATE column1 = VALUES(column1), column2 = VALUES(column2), ...; ステートメントが実行されます。 注:物理テーブルにプライマリキー以外の一意インデックス制約がある場合、異なるプライマリキーを持つが同じ一意インデックス値を持つ 2 つのレコードを挿入すると、一意インデックスの競合が発生します。 これにより、データが上書きされ、失われる可能性があります。

  • MySQL データベースで自動インクリメント主キーが定義されている場合、Flink DDL で自動インクリメントカラムを宣言しないでください。 データベースは、データ書き込み時にこのカラムに自動的に値を設定します。 コネクタは、自動インクリメントカラムを含むデータの書き込みと削除はサポートしますが、このデータの更新はサポートしていません。

WITH パラメータ

  • 一般

    パラメータ

    説明

    必須

    データ型

    デフォルト値

    備考

    connector

    テーブルタイプ。

    はい

    STRING

    なし

    ソーステーブルとして使用する場合、このパラメータは mysql-cdc または mysql に設定できます。これらは同等です。ディメンションテーブルまたはシンクテーブルとして使用する場合、値は mysql にする必要があります。

    hostname

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

    はい

    STRING

    なし

    Virtual Private Cloud (VPC) のアドレスを指定することを推奨します。

    説明

    MySQL データベースと Realtime Compute for Apache Flink が同じ Virtual Private Cloud (VPC) にない場合は、VPC 間ネットワーク接続を確立するか、パブリックエンドポイントを使用してデータベースにアクセスする必要があります。詳細については、「ワークスペースの管理と運用」および「フルマネージド Flink クラスターからインターネットにアクセスする方法」をご参照ください。

    username

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

    はい

    STRING

    なし

    なし。

    password

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

    はい

    STRING

    なし

    なし。

    database-name

    MySQL データベース名。

    はい

    STRING

    なし

    • データベースをソーステーブルとして使用する場合、データベース名に正規表現を使用して複数のデータベースからデータを読み取ることができます。

    • 正規表現を使用する場合、文字列の先頭と末尾に一致させるために ^ および $ 記号を使用しないでください。詳細については、table-name パラメータの備考をご参照ください。

    table-name

    MySQL テーブル名。

    はい

    STRING

    なし

    • ソーステーブル名に正規表現を使用して複数のテーブルからデータを読み取ることができます。

      複数の MySQL テーブルからデータを読み取る場合は、複数の CTAS ステートメントを 1 つのジョブとして送信してください。これにより、複数のバイナリログリスナーを有効にする必要がなくなり、パフォーマンスと効率が向上します。詳細については、「複数の CTAS ステートメント:単一ジョブとして送信」をご参照ください。

    • 正規表現を使用する場合、文字列の先頭と末尾に一致させるために ^ および $ 記号を使用しないでください。詳細については、次の注をご参照ください。

    説明

    MySQL CDC ソーステーブルが正規表現を使用してテーブル名を照合する場合、指定した database-nametable-name を文字列 \\. で連結し、フルパスの正規表現を作成します。VVR 8.0.1 より前のバージョンでは、文字 . が使用されていました。その後、コネクタはこの正規表現を使用して、MySQL データベース内のテーブルの完全修飾名を照合します。

    たとえば、'database-name'='db_.*' および 'table-name'='tb_.+' を設定した場合、コネクタは正規表現 db_.*\\.tb_.+ を使用して完全修飾テーブル名を照合し、読み取るテーブルを判定します。VVR 8.0.1 より前のバージョンでは、正規表現は db_.*.tb_.+ でした。

    port

    MySQL データベースサービスのポート番号。

    いいえ

    INTEGER

    3306

    なし。

  • ソーステーブルのみ

    パラメーター

    説明

    必須

    データ型

    デフォルト値

    備考

    server-id

    データベースクライアントの数値 ID を指定します。

    いいえ

    STRING

    5400~6400 の間でランダムな値が生成されます。

    この ID は、MySQL クラスター内でグローバルに一意である必要があります。同じデータベースに接続する各ジョブには、異なる ID を設定してください。

    このパラメーターは、5400-5408 などの ID 範囲形式もサポートしています。増分読み取りが有効な場合、同時読み取りがサポートされます。この場合、同時リーダーごとに異なる ID を使用するように ID 範囲を設定してください。詳細については、「Use Server ID」をご参照ください。

    scan.incremental.snapshot.enabled

    増分スナップショットを有効にするかどうかを指定します。

    いいえ

    BOOLEAN

    true

    増分スナップショットはデフォルトで有効です。増分スナップショットは、全データのスナップショットを読み取るための新しい仕組みです。従来のスナップショット読み取り方法と比べて、増分スナップショットには次のような多くの利点があります:

    • ソースは全データを並列に読み取ることができます。

    • ソースは、全データの読み取り時にチャンクレベルのチェックポイントをサポートします。

    • ソースは、全データの読み取り時にグローバル読み取りロック (FLUSH TABLES WITH read lock) を取得する必要がありません。

    ソースで同時読み取りをサポートする場合、同時リーダーごとに一意の server ID が必要です。そのため、server-id は 5400-6400 などの範囲である必要があり、範囲のサイズは並列度以上である必要があります。

    説明

    この設定項目は、Ververica Runtime (VVR) 11.1 以降で削除されています。

    scan.incremental.snapshot.chunk.size

    各チャンクのサイズ (行数) を指定します。

    いいえ

    INTEGER

    8096

    増分スナップショット読み取りが有効な場合、テーブルは複数のチャンクに分割されて読み取られます。チャンクのデータは、完全に読み取られるまでメモリにキャッシュされます。

    各チャンクに含まれる行数が少ないほど、テーブル内のチャンクの総数は多くなります。これにより障害復旧の粒度は細かくなりますが、メモリ不足 (OOM) エラーが発生する可能性があり、全体のスループットが低下する場合があります。そのため、トレードオフを考慮し、適切なチャンクサイズを設定する必要があります。

    scan.snapshot.fetch.size

    テーブルの全データを読み取る際に、1 回にプルするレコードの最大数を指定します。

    いいえ

    INTEGER

    1024

    なし。

    scan.startup.mode

    データ取り込みの起動モードを指定します。

    いいえ

    STRING

    initial

    有効な値:

    • initial (デフォルト):初回起動時に、コネクターはすべての履歴データをスキャンし、その後、最新のバイナリログデータを読み取ります。

    • latest-offset:初回起動時に、コネクターは履歴データをスキャンしません。バイナリログの末尾から読み取りを開始します。これは、コネクターの起動後に行われた最新の変更のみを読み取ることを意味します。

    • earliest-offset:コネクターは履歴データをスキャンしません。利用可能な最も古いバイナリログから読み取りを開始します。

    • specific-offset:コネクターは履歴データをスキャンしません。特定のバイナリログのオフセットから開始します。scan.startup.specific-offset.filescan.startup.specific-offset.pos の両方を設定するか、scan.startup.specific-offset.gtid-set のみを設定して特定の GTID セットから開始することで、オフセットを指定できます。

    • timestamp:コネクターは履歴データをスキャンしません。指定したタイムスタンプからバイナリログの読み取りを開始します。タイムスタンプは、scan.startup.timestamp-millis でミリ秒単位で指定します。

    重要

    earliest-offsetspecific-offset、または timestamp の起動モードを使用する場合は、指定したバイナリログの取り込み位置からジョブの起動時刻までの間に、対応するテーブルのスキーマが変更されていないことを確認してください。これにより、スキーマの不一致によるエラーを防止できます。

    scan.startup.specific-offset.file

    specific-offset 起動モードで使用する、開始オフセットのバイナリログファイル名を指定します。

    いいえ

    STRING

    なし

    このパラメーターを使用する場合、scan.startup.modespecific-offset に設定する必要があります。ファイル名形式の例: mysql-bin.000003

    scan.startup.specific-offset.pos

    specific-offset 起動モードで使用する、指定したバイナリログファイル内の開始オフセットを指定します。

    いいえ

    INTEGER

    なし

    このパラメーターを使用する場合は、scan.startup.modespecific-offset に設定する必要があります。

    scan.startup.specific-offset.gtid-set

    specific-offset 起動モードで使用する、開始オフセットの GTID セットを指定します。

    いいえ

    STRING

    なし

    このパラメーターを使用する場合、scan.startup.modespecific-offset に設定する必要があります。 GTID セット形式の例: 24DA167-0C0C-11E8-8442-00059A3C7B00:1-19

    scan.startup.timestamp-millis

    timestamp 起動モードで使用する、開始オフセットのタイムスタンプ (ミリ秒) を指定します。

    いいえ

    LONG

    なし

    このパラメーターを使用する場合は、scan.startup.modetimestamp に設定する必要があります。タイムスタンプの単位はミリ秒です。

    重要

    時刻を指定すると、MySQL CDC は各バイナリログファイルの初期イベントを読み取り、そのタイムスタンプを判定します。その後、指定した時刻に対応するバイナリログファイルを特定します。指定したタイムスタンプに対応するバイナリログファイルがデータベースから削除されておらず、読み取り可能であることを確認してください。

    server-time-zone

    データベースが使用するセッションタイムゾーンを指定します。

    いいえ

    STRING

    このパラメーターを指定しない場合、システムは Flink ジョブランタイムの環境タイムゾーンをデータベースサーバーのタイムゾーンとして使用します。これは、選択したゾーンのタイムゾーンです。

    例:Asia/Shanghai。このパラメーターは、MySQL の TIMESTAMP 型を STRING 型に変換する方法を制御します。詳細については、「Debezium temporal values」をご参照ください。

    debezium.min.row.count.to.stream.results

    テーブルの行数がこの値を超える場合、バッチ読み取りモードが使用されます。

    いいえ

    INTEGER

    1000

    Flink は、MySQL ソーステーブルから次のいずれかの方法でデータを読み取ります:

    • フル読み取り:テーブル全体のデータを直接メモリに読み取ります。この方法は高速ですが、それに応じたメモリを消費します。ソーステーブルが非常に大きい場合、OOM エラーが発生するリスクがあります。

    • バッチ読み取り:一定数の行を 1 バッチとして複数回に分けて読み取り、すべてのデータを読み取るまで繰り返します。この方法は、大きなテーブルの読み取り時に OOM リスクを回避できますが、比較的低速です。

    connect.timeout

    MySQL データベースサーバーへの接続がタイムアウトするまでの最大待機時間を指定します。タイムアウト後に再試行します。

    いいえ

    DURATION

    30s

    なし。

    connect.max-retries

    MySQL データベースサービスへの接続に失敗した後の最大再試行回数を指定します。

    いいえ

    INTEGER

    3

    なし。

    connection.pool.size

    データベース接続プールのサイズを指定します。

    いいえ

    INTEGER

    20

    データベース接続プールは接続を再利用するために使用され、データベース接続数を削減できます。

    jdbc.properties.*

    JDBC URL のカスタム接続パラメーターを指定します。

    いいえ

    STRING

    なし

    カスタム接続パラメーターを渡すことができます。たとえば、SSL プロトコルを使用しない場合は、'jdbc.properties.useSSL' = 'false' を設定できます。

    サポートされる接続パラメーターの詳細については、「MySQL Configuration Properties」をご参照ください。

    debezium.*

    Debezium がバイナリログを読み取るためのカスタムパラメーターを指定します。

    いいえ

    STRING

    なし

    カスタム Debezium パラメーターを渡すことができます。たとえば、デシリアライズエラーの処理ロジックを指定するには、'debezium.event.deserialization.failure.handling.mode'='ignore' を使用します。

    警告

    Debezium パラメーターを任意に変更しないでください。コネクターがデータを正しく読み取れなくなる可能性があります。たとえば、debezium.binlog.buffer.size パラメーターは設定できません。

    heartbeat.interval

    ハートビートイベントを使用して、ソースがバイナリログのオフセットを進める間隔を指定します。

    いいえ

    DURATION

    30s

    ハートビートイベントは、ソース内のバイナリログのオフセットを進めるために使用されます。これは、MySQL で更新頻度が低いテーブルで特に有用です。このようなテーブルでは、バイナリログのオフセットが自動的に進みません。ハートビートイベントによりバイナリログのオフセットを前進させ、バイナリログのオフセットの期限切れによって発生する問題を防止できます。バイナリログのオフセットが期限切れになると、ジョブが失敗して復旧できなくなり、ステートレス再起動が必要になる場合があります。

    scan.incremental.snapshot.chunk.key-column

    スナップショットフェーズ中のシャーディングにおける分割カラムとして使用するカラムを指定します。

    備考列をご参照ください。

    STRING

    なし

    • プライマリキーのないテーブルでは必須です。選択するカラムは非 NULL 型 (NOT NULL) である必要があります。

    • プライマリキーのあるテーブルでは任意です。プライマリキーから選択できるカラムは 1 つのみです。

    rds.region-id

    Alibaba Cloud ApsaraDB RDS for MySQL インスタンスのリージョン ID を指定します。

    OSS からアーカイブログを読み取る場合に必須です。

    STRING

    なし

    リージョン ID の詳細については、「Regions and zones」をご参照ください。

    重要

    MySQL CDC の GTID 文字列はランダムに生成され、バイナリログファイルのオフセットのように単調増加しないため、ファイル内の GTID を特定するには OSS からすべてのアーカイブログをダウンロードして解析する必要があります。この処理は大量のリソースを消費し、時間もかかるため、GTID オフセットに依存する機能は実現できません。そのため、OSS のアーカイブログ機能は、指定したタイムスタンプ、または指定したバイナリログファイルオフセットからの開始のみをサポートします。指定した GTID からの開始はサポートしていません。また、アーカイブログ内でのプライマリ/セカンダリ切り替えのシナリオもサポートしていません。MySQL のプライマリ/セカンダリ切り替えは GTID に依存しているためです。本機能は、使用前に十分に評価してください。

    rds.access-key-id

    Alibaba Cloud ApsaraDB RDS for MySQL アカウントの AccessKey ID を指定します。

    OSS からアーカイブログを読み取る場合に必須です。

    STRING

    なし

    詳細については、「How do I view the AccessKey ID and AccessKey secret?」をご参照ください。

    重要

    AccessKey 情報の漏えいを防ぐため、シークレット管理機能を使用して AccessKey ID を指定してください。詳細については、「Manage variables」をご参照ください。

    rds.access-key-secret

    Alibaba Cloud ApsaraDB RDS for MySQL アカウントの AccessKey Secret を指定します。

    OSS からアーカイブログを読み取る場合に必須です。

    STRING

    なし

    詳細については、「How do I view the AccessKey ID and AccessKey secret?」をご参照ください。

    重要

    AccessKey 情報の漏えいを防ぐため、シークレット管理機能を使用して AccessKey Secret を指定してください。詳細については、「Manage variables」をご参照ください。

    rds.db-instance-id

    Alibaba Cloud ApsaraDB RDS for MySQL インスタンスの ID を指定します。

    OSS からアーカイブログを読み取る場合に必須です。

    STRING

    なし

    なし。

    rds.main-db-id

    Alibaba Cloud ApsaraDB RDS for MySQL インスタンスのプライマリデータベース番号を指定します。

    いいえ

    STRING

    なし

    • プライマリデータベース番号の取得方法の詳細については、「ApsaraDB RDS for MySQL log backup」をご参照ください。

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

    説明

    このパラメーターを指定しない場合、VVR 11.7 以降では ApsaraDB RDS for MySQL の接続情報に基づいてプライマリデータベース番号が自動的に照会されます。

    rds.download.timeout

    OSS から単一のアーカイブログをダウンロードする際のタイムアウト期間を指定します。

    いいえ

    DURATION

    60s

    なし。

    rds.endpoint

    OSS のバイナリログ情報を取得するためのサービスエンドポイントを指定します。

    いいえ

    STRING

    なし

    • 有効な値の詳細については、「Endpoints」をご参照ください。

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

    scan.incremental.close-idle-reader.enabled

    スナップショットの完了後にアイドル状態のリーダーを閉じるかどうかを指定します。

    いいえ

    BOOLEAN

    false

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

    • この設定を有効にするには、execution.checkpointing.checkpoints-after-tasks-finish.enabled を true に設定する必要があります。

    scan.read-changelog-as-append-only.enabled

    変更ログデータストリームを追記専用データストリームに変換するかどうかを指定します。

    いいえ

    BOOLEAN

    false

    有効な値:

    • true:INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER を含むすべてのタイプのメッセージが INSERT メッセージに変換されます。アップストリームテーブルからの削除メッセージを保持する必要がある場合など、特別なシナリオでのみこのオプションを有効にしてください。

    • false (デフォルト):すべてのタイプのメッセージが、そのままダウンストリームに送信されます。

    説明

    VVR 8.0.8 以降でのみサポートされています。

    scan.only.deserialize.captured.tables.changelog.enabled

    増分フェーズで、指定したテーブルの変更イベントのみをデシリアライズするかどうかを指定します。

    いいえ

    BOOLEAN

    • VVR 8.x バージョンでは、デフォルト値は false です。

    • VVR 11.1 以降では、デフォルト値は true です。

    有効な値:

    • true:対象テーブルの変更データのみをデシリアライズし、バイナリログの読み取りを高速化します。

    • false:すべてのテーブルの変更データをデシリアライズします。

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

    • VVR 8.0.8 以前を使用する場合は、パラメーター名を debezium.scan.only.deserialize.captured.tables.changelog.enable に変更する必要があります。

    scan.parse.online.schema.changes.enabled

    増分フェーズで、RDS のロックレス変更 DDL イベントの解析を試行するかどうかを指定します。

    いいえ

    BOOLEAN

    false

    有効な値:

    • true:RDS のロックレス変更 DDL イベントを解析します。

    • false (デフォルト):RDS のロックレス変更 DDL イベントを解析しません。

    これは実験的な機能です。オンラインでロックレス変更を実行する前に、復旧用に Flink ジョブのスナップショットを取得してください。

    説明

    VVR 11.1 以降でのみサポートされています。

    scan.incremental.snapshot.backfill.skip

    スナップショット読み取りフェーズ中にバックフィルをスキップするかどうかを指定します。

    いいえ

    BOOLEAN

    false

    有効な値:

    • true:スナップショット読み取りフェーズ中のバックフィルをスキップします。

    • false (デフォルト):スナップショット読み取りフェーズ中のバックフィルをスキップしません。

    バックフィルをスキップすると、スナップショットフェーズ中に発生したテーブルの変更は、スナップショットにマージされず、後続の増分フェーズで読み取られます。

    重要

    バックフィルをスキップすると、スナップショットフェーズ中に発生した変更が再生される可能性があるため、データの不整合が発生する可能性があります。at-least-once セマンティクスのみが保証されます。

    説明

    VVR 11.1 以降でのみサポートされています。

    scan.incremental.snapshot.unbounded-chunk-first.enabled

    スナップショット読み取りフェーズ中に、無制限チャンクを先にディスパッチするかどうかを指定します。

    いいえ

    BOOLEAN

    false

    有効な値:

    • true:スナップショット読み取りフェーズ中に、無制限チャンクを先にディスパッチします。

    • false (デフォルト):スナップショット読み取りフェーズ中に、無制限チャンクを先にディスパッチしません。

    これは実験的な機能です。有効にすると、スナップショットフェーズ中に最後のチャンクを同期する際の TaskManager でのメモリ不足 (OOM) エラーのリスクを低減できます。ジョブの初回起動前にこのパラメーターを追加してください。

    説明

    VVR 11.1 以降でのみサポートされています。

    binlog.session.network.timeout

    バイナリログ接続のネットワーク読み取り/書き込みタイムアウトを指定します。

    いいえ

    DURATION

    10m

    0s に設定した場合、MySQL サーバーのデフォルトタイムアウトが使用されます。

    説明

    VVR 11.5 以降でのみサポートされています。

    scan.rate-limit.records-per-second

    ソースが 1 秒あたりに送信するレコードの最大数を制限します。

    いいえ

    LONG

    なし

    データ読み取りを制限する必要があるシナリオで使用できます。この制限は、全フェーズと増分フェーズの両方で有効です。

    ソースの numRecordsOutPerSecond メトリクスは、データストリーム全体が 1 秒あたりに出力するレコード数を反映します。このメトリクスに基づいて本パラメーターを調整できます。

    全データ読み取りフェーズでは、通常、各バッチで読み取る行数を減らす必要があります。scan.incremental.snapshot.chunk.size パラメーターの値を減らすことができます。

    説明

    VVR 11.5 以降でのみサポートされています。

    scan.binlog.tolerate.gtid-holes

    このパラメーターを有効にすると、GTID の順序の欠落が無視され、ジョブは不連続なイベントをスキップして実行を継続できます。

    いいえ

    BOOLEAN

    false

    このパラメーターを有効にする前に、ジョブの開始オフセットが期限切れになっていないことを確認する必要があります。ジョブが削除済みまたは期限切れの GTID オフセットから開始すると、エンジンは欠落したログを黙ってスキップし、データ損失が発生します。

    説明

    このパラメーターは VVR 11.6 以降でのみサポートされています。

  • ディメンションテーブル固有のパラメーター

    パラメーター

    説明

    必須

    データ型

    デフォルト値

    備考

    url

    MySQL の JDBC URL です。

    いいえ

    STRING

    なし

    URL の形式は jdbc:mysql://<endpoint>:<port>/<database_name> です。

    lookup.max-retries

    データ読み取りに失敗した後の最大リトライ回数です。

    いいえ

    INTEGER

    3

    VVR 6.0.7 以降でのみサポートされています。

    lookup.cache.strategy

    キャッシュポリシーです。

    いいえ

    STRING

    None

    サポートされているキャッシュポリシーは None、LRU、ALL です。 値の詳細については、「ディメンションテーブルの JOIN ステートメント」をご参照ください。

    説明

    LRU キャッシュポリシーを使用する場合は、lookup.cache.max-rows パラメーターも設定する必要があります。

    lookup.cache.max-rows

    キャッシュされた行の最大数です。

    いいえ

    INTEGER

    100000

    • LRU キャッシュポリシーを選択した場合は、キャッシュサイズを設定する必要があります。

    • ALL キャッシュポリシーを選択した場合は、キャッシュサイズを設定する必要はありません。

    lookup.cache.ttl

    キャッシュの Time-to-Live (TTL) です。

    いいえ

    DURATION

    なし

    lookup.cache.ttl の設定は lookup.cache.strategy によって異なります。

    • lookup.cache.strategyNone に設定されている場合、lookup.cache.ttl の設定は不要です。この場合、キャッシュはタイムアウトしません。

    • lookup.cache.strategyLRU に設定されている場合、lookup.cache.ttl はキャッシュの TTL です。 デフォルトでは、キャッシュは期限切れになりません。

    • lookup.cache.strategyALL に設定されている場合、lookup.cache.ttl はキャッシュの読み込み時間です。 デフォルトでは、キャッシュは再読み込みされません。

    1min や 10s などの時間形式を使用します。

    lookup.max-join-rows

    プライマリテーブルのレコードがディメンションテーブルのレコードと一致した場合に返される、結果の最大数です。

    いいえ

    INTEGER

    1024

    なし

    lookup.filter-push-down.enabled

    ディメンションテーブルのフィルタープッシュダウンを有効にするかどうかを指定します。

    いいえ

    BOOLEAN

    false

    有効な値:

    • true:ディメンションテーブルのフィルタープッシュダウンを有効にします。 MySQL データベーステーブルからデータをロードする際、ディメンションテーブルは SQL ジョブで設定された条件に基づき、事前にデータをフィルターします。

    • false (デフォルト):ディメンションテーブルのフィルタープッシュダウンを無効にします。 MySQL データベーステーブルからデータをロードする際、ディメンションテーブルはすべてのデータをロードします。

    説明

    VVR 8.0.7 以降でのみサポートされています。

    重要

    ディメンションテーブルのプッシュダウンは、Flink テーブルがディメンションテーブルとして使用される場合にのみ有効にしてください。 MySQL ソーステーブルは、フィルタープッシュダウンの有効化をサポートしていません。 Flink テーブルをソーステーブルとディメンションテーブルの両方として使用し、かつディメンションテーブルでフィルタープッシュダウンを有効にする場合、SQL ヒントを使用して、ソーステーブルに対してこの設定項目を明示的に false に設定する必要があります。 そうしないと、ジョブが異常に実行される可能性があります。

  • シンクテーブルのみ

    パラメーター

    説明

    必須

    データ型

    デフォルト値

    備考

    url

    MySQL JDBC URL。

    任意

    STRING

    なし

    URL 形式: jdbc:mysql://<endpoint>:<port>/<database_name>

    sink.max-retries

    データ書き込み失敗時の最大リトライ回数。

    任意

    INTEGER

    3

    なし。

    sink.buffer-flush.batch-size

    1 回のバッチ書き込みにおける行数。

    任意

    INTEGER

    4096

    なし。

    sink.buffer-flush.max-rows

    メモリにキャッシュされるデータ行数。

    任意

    INTEGER

    10000

    このパラメーターは、プライマリキーを指定した場合にのみ有効です。

    sink.buffer-flush.interval

    キャッシュをフラッシュする間隔。指定した待機時間が経過してもキャッシュ内のデータが出力条件を満たさない場合、システムはキャッシュ内のすべてのデータを自動的に出力します。

    任意

    DURATION

    1 s

    なし。

    sink.ignore-delete

    データの DELETE 操作を無視するかどうかを指定します。

    任意

    BOOLEAN

    false

    Flink SQL で生成されるストリームに delete または update-before レコードが含まれる場合、複数の出力タスクが同じテーブルの異なるフィールドを同時に更新すると、データの不整合が発生する可能性があります。

    たとえば、レコードが削除された後に、別のタスクが一部のフィールドのみを更新した場合、更新されなかったフィールドが null またはデフォルト値となり、データエラーが発生します。

    sink.ignore-delete を true に設定すると、アップストリームの DELETE 操作と UPDATE_BEFORE 操作を無視して、このような問題を回避できます。

    説明
    • UPDATE_BEFORE は Flink のリトラクションメカニズムの一部であり、更新操作における古い値を「取り消す」ために使用されます。

    • ignoreDelete = true の場合、すべての DELETE および UPDATE_BEFORE タイプのレコードがスキップされます。INSERT および UPDATE_AFTER レコードのみが処理されます。

    sink.ignore-null-when-update

    データを更新する際、入力データのフィールド値が null であれば、対応するフィールドを null に更新するか、そのフィールドの更新をスキップするかを指定します。

    任意

    BOOLEAN

    false

    有効な値:

    • true:フィールドを更新しません。このパラメーターは、Flink テーブルにプライマリキーが設定されている場合にのみ true に設定できます。true に設定した場合:

      • VVR 8.0.6 以前では、シンクテーブルはバッチ書き込みをサポートしていません。

      • VVR 8.0.7 以降では、シンクテーブルはバッチ書き込みをサポートしています。

        バッチ書き込みは書き込み効率と全体のスループットを大幅に向上させますが、データレイテンシーと OOM エラーのリスクを伴います。そのため、ビジネスシナリオに応じてトレードオフを検討する必要があります。

    • false:フィールドを null に更新します。

    説明

    このパラメーターは、VVR 8.0.5 以降でのみサポートされます。

型マッピング

  • CDC ソーステーブル

    MySQL CDC フィールドタイプ

    Flink フィールドタイプ

    TINYINT

    TINYINT

    SMALLINT

    SMALLINT

    TINYINT UNSIGNED

    TINYINT UNSIGNED ZEROFILL

    INT

    INT

    MEDIUMINT

    SMALLINT UNSIGNED

    SMALLINT UNSIGNED ZEROFILL

    BIGINT

    BIGINT

    INT UNSIGNED

    INT UNSIGNED ZEROFILL

    MEDIUMINT UNSIGNED

    MEDIUMINT UNSIGNED ZEROFILL

    BIGINT UNSIGNED

    DECIMAL(20, 0)

    BIGINT UNSIGNED ZEROFILL

    SERIAL

    FLOAT [UNSIGNED] [ZEROFILL]

    FLOAT

    DOUBLE [UNSIGNED] [ZEROFILL]

    DOUBLE

    DOUBLE PRECISION [UNSIGNED] [ZEROFILL]

    REAL [UNSIGNED] [ZEROFILL]

    NUMERIC(p, s) [UNSIGNED] [ZEROFILL]

    DECIMAL(p, s)

    DECIMAL(p, s) [UNSIGNED] [ZEROFILL]

    BOOLEAN

    BOOLEAN

    TINYINT(1)

    DATE

    DATE

    TIME [(p)]

    TIME [(p)] [WITHOUT TIME ZONE]

    DATETIME [(p)]

    TIMESTAMP [(p)] [WITHOUT TIME ZONE]

    TIMESTAMP [(p)]

    TIMESTAMP [(p)]

    TIMESTAMP(p) WITH LOCAL TIME ZONE

    CHAR(n)

    STRING

    VARCHAR(n)

    TEXT

    BINARY

    BYTES

    VARBINARY

    BLOB

    重要

    MySQL で TINYINT(1) 型を使用して 0 と 1 以外の値を格納しないでください。property-version=0 の場合、MySQL CDC ソーステーブルはデフォルトで TINYINT(1) を Flink の BOOLEAN 型にマップします。これにより、データが不正確になる可能性があります。TINYINT(1) 型を使用して 0 と 1 以外の値を格納する場合は、設定パラメーター catalog.table.treat-tinyint1-as-boolean をご参照ください。

  • ディメンションテーブルおよびシンクテーブル

    MySQL フィールドタイプ

    Flink フィールドタイプ

    TINYINT

    TINYINT

    SMALLINT

    SMALLINT

    TINYINT UNSIGNED

    INT

    INT

    MEDIUMINT

    SMALLINT UNSIGNED

    BIGINT

    BIGINT

    INT UNSIGNED

    BIGINT UNSIGNED

    DECIMAL(20, 0)

    FLOAT

    FLOAT

    DOUBLE

    DOUBLE

    DOUBLE PRECISION

    NUMERIC(p, s)

    DECIMAL(p, s)

    説明

    ただし、p <= 38。

    DECIMAL(p, s)

    BOOLEAN

    BOOLEAN

    TINYINT(1)

    DATE

    DATE

    TIME [(p)]

    TIME [(p)] [WITHOUT TIME ZONE]

    DATETIME [(p)]

    TIMESTAMP [(p)] [WITHOUT TIME ZONE]

    TIMESTAMP [(p)]

    CHAR(n)

    CHAR(n)

    VARCHAR(n)

    VARCHAR(n)

    BIT(n)

    BINARY(⌈n/8⌉)

    BINARY(n)

    BINARY(n)

    VARBINARY(N)

    VARBINARY(N)

    TINYTEXT

    STRING

    TEXT

    MEDIUMTEXT

    LONGTEXT

    TINYBLOB

    BYTES

    重要

    Flink は、2,147,483,647 (2^31 - 1) バイト以下の MySQL BLOB 型レコードのみをサポートします。

    BLOB

    MEDIUMBLOB

    LONGBLOB

データインジェスト

データインジェストの YAML ジョブでは、MySQL コネクタをデータソースとして使用できます。

構文

source:
   type: mysql
   name: MySQL Source
   hostname: localhost
   port: 3306
   username: <username>
   password: <password>
   tables: adb.\.*, bdb.user_table_[0-9]+, [app|web].order_\.*
   server-id: 5401-5404

sink:
  type: xxx

構成アイテム

パラメータ

説明

必須

データ型

デフォルト値

備考

type

データソースのタイプ。

はい

STRING

なし

値はmysqlである必要があります。

name

データソースの名前。

いいえ

STRING

なし

なし。

hostname

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

はい

STRING

なし

VPC アドレスを指定することを推奨します。

説明

MySQL データベースと Realtime Compute for Apache Flink が同じ VPC にない場合は、VPC 間ネットワーク接続を確立するか、パブリックエンドポイントを使用してデータベースにアクセスする必要があります。 詳細については、「ワークスペースの管理と運用」および「フルマネージド Flink クラスターはどのようにインターネットにアクセスしますか?」をご参照ください。

username

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

はい

STRING

なし

なし。

password

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

はい

STRING

なし

なし。

tables

同期する MySQL のデータテーブル。

はい

STRING

なし

  • このパラメータは正規表現に対応しており、複数のテーブルからデータを読み取れます。

  • カンマで複数の正規表現を区切ります。

説明
  • 正規表現では、文字列の先頭を表す^や末尾を表す$のマッチング文字を使用しないでください。 VVR 11.2 では、ピリオドを使用して正規表現を分割し、データベース部分を取得します。 これらの文字を使用すると、分割後のデータベース正規表現が使用できなくなります。 たとえば、^db.user_[0-9]+$db.user_[0-9]+に変更する必要があります。

  • ピリオドはデータベース名とテーブル名を区切るために使用します。 ピリオドを任意の文字にマッチさせるには、バックスラッシュでエスケープする必要があります。 例:db0.\.*db1.user_table_[0-9]+db[1-2].[app|web]order_\.*

tables.exclude

同期対象から除外するテーブル。

いいえ

STRING

なし

  • このパラメータは正規表現に対応しており、複数のテーブルを除外できます。

  • カンマで複数の正規表現を区切ります。

説明

ピリオドはデータベース名とテーブル名を区切るために使用します。 ピリオドを任意の文字にマッチさせるには、バックスラッシュでエスケープする必要があります。 例:db0.\.*db1.user_table_[0-9]+db[1-2].[app|web]order_\.*

port

MySQL データベースサービスのポート番号。

いいえ

INTEGER

3306

なし。

schema-change.enabled

スキーマ変更イベントを送信するかどうかを指定します。

いいえ

BOOLEAN

true

なし。

server-id

同期に使用するデータベースクライアントの数値 ID または範囲。

いいえ

STRING

5400~6400 のランダムな値が生成されます。

この ID は MySQL クラスター内で一意である必要があります。 同じデータベースに接続する各ジョブには、異なる ID を設定します。

このパラメータは、5400-5408などの ID 範囲形式もサポートしています。 増分読み取りが有効な場合、同時読み取りが可能です。 この場合、各同時読み取りリーダーで異なる ID が使用されるように、ID 範囲を設定します。

jdbc.properties.*

JDBC URL のカスタム接続パラメータ。

いいえ

STRING

なし

カスタム接続パラメータを指定できます。 たとえば、SSL プロトコルを使用しない場合は、'jdbc.properties.useSSL' = 'false'を設定できます。

利用可能な接続パラメータの詳細については、「MySQL 設定プロパティ」をご参照ください。

debezium.*

Debezium がバイナリログを読み取るためのカスタムパラメータ。

いいえ

STRING

なし

カスタムの Debezium パラメータを指定できます。 たとえば、'debezium.event.deserialization.failure.handling.mode'='ignore'を使用して、パースエラーの処理ロジックを指定します。

警告

Debezium パラメータを任意に変更しないでください。コネクタがデータを誤って読み取る原因となります。 たとえば、debezium.binlog.buffer.sizeパラメータの設定はできません。

scan.incremental.snapshot.chunk.size

各チャンクのサイズ (行数)。

いいえ

INTEGER

8096

MySQL テーブルは、読み取り時に複数のチャンクに分割されます。 チャンクのデータは、完全に読み取られるまでメモリにキャッシュされます。

チャンクあたりの行数が少ないほど、テーブルのチャンク総数は多くなります。 この場合、障害回復の粒度は細かくなりますが、OOM エラーや全体的なスループットの低下につながる可能性があります。 したがって、トレードオフを考慮し、適切なチャンクサイズを設定する必要があります。

scan.snapshot.fetch.size

テーブルの全量データを読み取る際に一度に取得するレコードの最大数。

いいえ

INTEGER

1024

なし。

scan.startup.mode

データ消費の起動モード。

いいえ

STRING

initial

有効な値:

  • initial (デフォルト):初回起動時、コネクタは全量履歴データをスキャンし、最新のバイナリログデータを読み取ります。

  • latest-offset:初回起動時、コネクタは履歴データをスキャンしません。 バイナリログの末尾から読み取りを開始します。つまり、コネクタの起動後に行われた最新の変更のみを読み取ります。

  • earliest-offset:コネクタは履歴データをスキャンしません。 利用可能な最も古いバイナリログから読み取りを開始します。

  • specific-offset:コネクタは履歴データをスキャンしません。 特定のバイナリログオフセットから開始します。 scan.startup.specific-offset.filescan.startup.specific-offset.posの両方を設定するか、scan.startup.specific-offset.gtid-setのみ設定して特定の GTID セットから開始することで、オフセットを指定できます。

  • timestamp:コネクタは履歴データをスキャンしません。 指定されたタイムスタンプからバイナリログの読み取りを開始します。 タイムスタンプはscan.startup.timestamp-millisでミリ秒単位で指定します。

重要

earliest-offsetspecific-offsettimestampの起動モードでは、起動時のテーブルスキーマが指定された開始オフセット時刻のスキーマと異なる場合、スキーマの不一致によりジョブでエラーが発生します。 言い換えれば、これら 3 つの起動モードを使用する場合、指定されたバイナリログの読み取り開始位置とジョブの起動時刻の間で、対応するテーブルのスキーマが変更されていないことを確認してください。

scan.startup.specific-offset.file

specific-offset起動モードを使用する場合の開始オフセットのバイナリログファイル名。

いいえ

STRING

なし

このパラメータを使用する場合、scan.startup.modespecific-offsetに設定する必要があります。 ファイル名の例:mysql-bin.000003

scan.startup.specific-offset.pos

specific-offset起動モードを使用する場合の、指定されたバイナリログファイル内の開始オフセット。

いいえ

INTEGER

なし

このパラメータを使用する場合、scan.startup.modespecific-offsetに設定する必要があります。

scan.startup.specific-offset.gtid-set

specific-offset起動モードを使用する場合の開始オフセットの GTID セット。

いいえ

STRING

なし

このパラメータを使用する場合、scan.startup.modespecific-offsetに設定する必要があります。 GTID セットの例:24DA167-0C0C-11E8-8442-00059A3C7B00:1-19

scan.startup.timestamp-millis

timestamp起動モードを使用する場合の開始オフセットのタイムスタンプ (ミリ秒)。

いいえ

LONG

なし

このパラメータを使用する場合、scan.startup.modetimestampに設定する必要があります。 タイムスタンプの単位はミリ秒です。

重要

時刻を指定すると、MySQL CDC は各バイナリログファイルの初期イベントを読み取ってそのタイムスタンプを特定します。 次に、指定された時刻に対応するバイナリログファイルを見つけます。 指定されたタイムスタンプに対応するバイナリログファイルがデータベースから削除されておらず、読み取り可能であることを確認してください。

server-time-zone

データベースが使用するセッションタイムゾーン。

いいえ

STRING

このパラメータを指定しない場合、システムは Flink ジョブランタイムの環境タイムゾーンをデータベースサーバーのタイムゾーンとして使用します。 これは、選択したゾーンのタイムゾーンです。

例:Asia/Shanghai。 このパラメータは、MySQL の TIMESTAMP 型が STRING 型に変換される方法を制御します。 詳細については、「Debezium の時間的な値」をご参照ください。

scan.startup.specific-offset.skip-events

指定されたオフセットから読み取る際にスキップするバイナリログイベントの数。

いいえ

INTEGER

なし

このパラメータを使用する場合、scan.startup.modespecific-offsetに設定する必要があります。

scan.startup.specific-offset.skip-rows

指定されたオフセットから読み取る際にスキップする行変更の数。 単一のバイナリログイベントは、複数の行変更に対応する場合があります。

いいえ

INTEGER

なし

このパラメータを使用する場合、scan.startup.modespecific-offsetに設定する必要があります。

connect.timeout

MySQL データベースサーバーへの接続がタイムアウトし、再試行するまでの最大待機時間。

いいえ

DURATION

30 s

なし。

connect.max-retries

MySQL データベースサービスへの接続に失敗した後の最大再試行回数。

いいえ

INTEGER

3

なし。

connection.pool.size

データベース接続プールのサイズ。

いいえ

INTEGER

20

データベース接続プールは接続を再利用するために使用され、これによりデータベース接続の数を減らすことができます。

heartbeat.interval

ソースがハートビートイベントを使用してバイナリログのオフセットを進める間隔。

いいえ

DURATION

30s

ハートビートイベントは、ソース内のバイナリログオフセットを進めるために使用します。 これは、更新頻度の低い MySQL テーブルに非常に役立ちます。 このようなテーブルでは、バイナリログのオフセットは自動的に進みません。 ハートビートイベントはバイナリログのオフセットを進めることができるため、バイナリログオフセットの期限切れが原因で発生する問題を回避できます。 バイナリログオフセットが期限切れになると、ジョブが失敗して回復不能になり、ステートレスな再起動が必要になる場合があります。

rds.region-id

Alibaba Cloud ApsaraDB RDS for MySQL インスタンスのリージョン ID。

OSS からアーカイブログを読み取る機能を使用する場合は必須です。

STRING

なし

リージョン ID の詳細については、「リージョンとゾーン」をご参照ください。

重要

MySQL CDC の GTID 文字列はランダムに生成され、バイナリログファイルのオフセットのように単調に増加しないため、ファイル内で GTID を見つけるには、OSS からすべてのアーカイブログをダウンロードしてパースする必要があります。 このプロセスは非常にリソースを消費し、時間がかかるため、GTID オフセットに依存する機能は現実的ではありません。 したがって、OSS アーカイブログ機能は、指定されたタイムスタンプまたは指定されたバイナリログファイルオフセットからの開始のみをサポートします。 指定された GTID からの開始はサポートしていません。また、MySQL のプライマリ/セカンダリ切り替えは GTID に依存するため、アーカイブログでのプライマリ/セカンダリ切り替えシナリオもサポートしていません。 使用する前に、この機能を慎重に評価してください。

rds.access-key-id

Alibaba Cloud ApsaraDB RDS for MySQL アカウントの AccessKey ID。

OSS からアーカイブログを読み取る機能を使用する場合は必須です。

STRING

なし

詳細については、「AccessKey ID と AccessKey Secret を表示するにはどうすればよいですか?」をご参照ください。

重要

AccessKey 情報の漏洩を防ぐために、シークレット管理機能を使用して AccessKey ID を指定することを推奨します。 詳細については、「変数の管理」をご参照ください。

rds.access-key-secret

Alibaba Cloud ApsaraDB RDS for MySQL アカウントの AccessKey Secret。

OSS からアーカイブログを読み取る機能を使用する場合は必須です。

STRING

なし

詳細については、「AccessKey ID と AccessKey Secret を表示するにはどうすればよいですか?」をご参照ください。

重要

AccessKey 情報の漏洩を防ぐために、シークレット管理機能を使用して AccessKey Secret を指定することを推奨します。 詳細については、「変数の管理」をご参照ください。

rds.db-instance-id

Alibaba Cloud ApsaraDB RDS for MySQL インスタンスの ID。

OSS からアーカイブログを読み取る機能を使用する場合は必須です。

STRING

なし

なし。

rds.main-db-id

Alibaba Cloud ApsaraDB RDS for MySQL インスタンスのプライマリデータベース番号。

いいえ

STRING

なし

プライマリデータベース番号の取得方法の詳細については、「ApsaraDB RDS for MySQL のログバックアップ」をご参照ください。

説明

このパラメータが指定されていない場合、VVR 11.7 以降では、ApsaraDB RDS for MySQL 接続情報に基づいてプライマリデータベース番号が自動的に照会されます。

rds.download.timeout

OSS から単一のアーカイブログをダウンロードするためのタイムアウト期間。

いいえ

DURATION

60s

なし。

rds.endpoint

OSS バイナリログ情報を取得するためのサービスエンドポイント。

いいえ

STRING

なし

有効な値の詳細については、「エンドポイント」をご参照ください。

rds.binlog-directory-prefix

バイナリログファイルを格納するためのディレクトリプレフィックス。

いいえ

STRING

rds-binlog-

なし。

rds.use-intranet-link

内部ネットワークを使用してバイナリログファイルをダウンロードするかどうかを指定します。

いいえ

BOOLEAN

true

なし。

rds.binlog-directories-parent-path

バイナリログファイルを格納するための親ディレクトリの絶対パス。

いいえ

STRING

なし

なし。

chunk-meta.group.size

チャンクメタデータのサイズ。

いいえ

INTEGER

1000

メタデータがこの値より大きい場合、送信のために複数の部分に分割されます。

chunk-key.even-distribution.factor.lower-bound

均等シャーディングのためのチャンク分布係数の下限。

いいえ

DOUBLE

0.05

分布係数がこの値より小さい場合、不均等シャーディングが使用されます。

チャンク分布係数 = (MAX(chunk-key) - MIN(chunk-key) + 1) / データ行の総数。

chunk-key.even-distribution.factor.upper-bound

均等シャーディングのためのチャンク分布係数の上限。

いいえ

DOUBLE

1000.0

分布係数がこの値より大きい場合、不均等シャーディングが使用されます。

チャンク分布係数 = (MAX(chunk-key) - MIN(chunk-key) + 1) / データ行の総数。

scan.incremental.close-idle-reader.enabled

スナップショットが完了した後にアイドル状態のリーダーを閉じるかどうかを指定します。

いいえ

BOOLEAN

false

この設定を有効にするには、execution.checkpointing.checkpoints-after-tasks-finish.enabledtrueに設定する必要があります。

scan.only.deserialize.captured.tables.changelog.enabled

増分フェーズで、指定されたテーブルの変更イベントのみをデシリアライズするかどうかを指定します。

いいえ

BOOLEAN

  • VVR 8.x バージョンではデフォルト値はfalseです。

  • VVR 11.1 以降ではデフォルト値はtrueです。

有効な値:

  • true:対象テーブルの変更データのみをデシリアライズして、バイナリログの読み取りを高速化します。

  • false:すべてのテーブルの変更データをデシリアライズします。

scan.parallel-deserialize-changelog.enabled

増分フェーズで、複数のスレッドを使用して変更イベントをパースするかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:変更イベントのデシリアライズフェーズで複数のスレッドを使用し、バイナリログイベントの順序を維持しながら読み取りを高速化します。

  • false (デフォルト):イベントのデシリアライズフェーズで単一のスレッドを使用します。

説明

VVR 8.0.11 以降でのみサポートされています。

scan.parallel-deserialize-changelog.handler.size

複数のスレッドを使用して変更イベントをパースする際のイベントハンドラーの数。

いいえ

INTEGER

2

説明

VVR 8.0.11 以降でのみサポートされています。

metadata-column.include-list

ダウンストリームに渡すメタデータ列。

いいえ

STRING

なし

利用可能なメタデータには、op_tses_tsquery_logfileposが含まれます。 カンマで複数のメタデータ列を区切ることができます。

説明

MySQL CDC YAML コネクタは、データベース名、テーブル名、およびop_typeメタデータ列の追加を要求またはサポートしていません。 Transform 式で__data_event_type__を直接使用して変更データタイプを取得したり、__schema_name____table_name__を使用してデータベース名とテーブル名を取得したりできます。

重要
  • fileメタデータ列は、データが配置されているバイナリログファイルを表します。 全量フェーズでは""であり、増分フェーズではバイナリログファイル名です。 posメタデータ列は、バイナリログファイル内のデータのオフセットを表します。 全量フェーズでは"0"であり、増分フェーズではバイナリログファイル内のデータオフセットです。 これら 2 つのメタデータ列は、VVR 11.5 以降でサポートされています。

  • es_tsメタデータ列は、MySQL 上の変更ログに対応するトランザクションの開始時刻を表します。 MySQL 8.0.x でのみサポートされています。 以前のバージョンの MySQL を使用する場合は、このメタデータ列を追加しないでください。

  • op_tsタイムスタンプは秒単位で正確ですが、es_tsタイムスタンプはミリ秒単位で正確です。

scan.newly-added-table.enabled

チェックポイントから再起動する際に、前回の起動時にマッチしなかった新しく追加されたテーブルを同期するか、マッチしなくなったテーブルを状態から削除するかを指定します。

いいえ

BOOLEAN

false

これは、チェックポイントまたはセーブポイントから再起動するときに有効になります。

重要

全量データ読み取りフェーズ中に、セーブポイントを保存したり、ソーステーブルに新しいテーブルを追加または削除したり、その後セーブポイントからジョブを再起動したりすることはできません。 これにより、ジョブはデータの読み取りに失敗します。

scan.binlog.newly-added-table.enabled

増分フェーズで、マッチした新しく追加されたテーブルからデータを送信するかどうかを指定します。

いいえ

BOOLEAN

false

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

scan.incremental.snapshot.chunk.key-column

スナップショットフェーズ中にシャーディングの分割列として使用する、特定のテーブルの列を指定します。

いいえ

STRING

なし

  • コロン:を使用してテーブル名と列名を接続し、ルールを定義します。 テーブル名は正規表現にすることができます。 セミコロン;で区切ることで、複数のルールを定義できます。 例:db1.user_table_[0-9]+:col1;db[1-2].[app|web]_order_\\.*:col2

  • プライマリキーのないテーブルでは必須です。 選択された列は、非ヌル型 (NOT NULL) である必要があります。 プライマリキーのあるテーブルではオプションです。 プライマリキーから選択できる列は 1 つだけです。

scan.parse.online.schema.changes.enabled

増分フェーズで、RDS のロックレス変更 DDL イベントのパースを試みるかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:RDS のロックレス変更 DDL イベントをパースします。

  • false (デフォルト):RDS のロックレス変更 DDL イベントをパースしません。

これは実験的な機能です。 オンラインでのロックレス変更を実行する前に、Flink ジョブのスナップショットを取得して回復に備えてください。

説明

VVR 11.0 以降でのみサポートされています。

scan.incremental.snapshot.backfill.skip

スナップショット読み取りフェーズ中にバックフィルをスキップするかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:スナップショット読み取りフェーズ中にバックフィルをスキップします。

  • false (デフォルト):スナップショット読み取りフェーズ中にバックフィルをスキップしません。

バックフィルがスキップされた場合、スナップショットフェーズ中のテーブルへの変更は、スナップショットにマージされるのではなく、後の増分フェーズで読み取られます。

重要

バックフィルをスキップすると、スナップショットフェーズ中に発生した変更が再実行される可能性があるため、データの不整合につながる可能性があります。 at-least-once セマンティクスのみが保証されます。

説明

VVR 11.1 以降でのみサポートされています。

treat-tinyint1-as-boolean.enabled

TINYINT(1)型をブール型として扱うかどうかを指定します。

いいえ

BOOLEAN

true

有効な値:

  • true (デフォルト):TINYINT(1)型をブール型として扱います。

  • falseTINYINT(1)型をブール型として扱いません。

treat-timestamp-as-datetime-enabled

TIMESTAMP型をDATETIME型として扱うかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:MySQL のTIMESTAMP型をDATETIME型として扱い、CDC のTIMESTAMP型にマッピングします。

  • false (デフォルト):MySQL のTIMESTAMP型を CDC のTIMESTAMP_LTZ型にマッピングします。

MySQL のTIMESTAMP型は UTC 時刻を格納し、タイムゾーンの影響を受けます。 MySQL のDATETIME型はリテラルな時刻を格納し、タイムゾーンの影響を受けません。

有効にすると、server-time-zoneに基づいて MySQL のTIMESTAMP型データがDATETIME型に変換されます。

include-comments.enabled

テーブルと列のコメントを同期するかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:テーブルと列のコメントを同期します。

  • false (デフォルト):テーブルと列のコメントを同期しません。

このオプションを有効にすると、ジョブのメモリ使用量が増加します。

scan.incremental.snapshot.unbounded-chunk-first.enabled

スナップショット読み取りフェーズ中に、境界のないチャンクを最初にディスパッチするかどうかを指定します。

いいえ

BOOLEAN

false

有効な値:

  • true:スナップショット読み取りフェーズ中に、境界のないチャンクを最初にディスパッチします。

  • false (デフォルト):スナップショット読み取りフェーズ中に、境界のないチャンクを最初にディスパッチしません。

これは実験的な機能です。 有効にすると、スナップショットフェーズ中に最後のチャンクを同期する際の TaskManager での OOM エラーのリスクを軽減できます。 ジョブの初回起動前にこのパラメータを追加してください。

説明

VVR 11.1 以降でのみサポートされています。

binlog.session.network.timeout

バイナリログ接続のネットワークタイムアウト。

いいえ

DURATION

10m

0sに設定した場合、MySQL サーバーのデフォルトのタイムアウトが使用されます。

説明

VVR 11.5 以降でのみサポートされています。

scan.rate-limit.records-per-second

ソースが 1 秒あたりに送信するレコードの最大数を制限します。

いいえ

LONG

なし

これは、データ読み取りを制限する必要があるシナリオに適用できます。 この制限は、全量フェーズと増分フェーズの両方で有効です。

ソースの numRecordsOutPerSecond メトリックは、データストリーム全体から毎秒出力されるレコード数を示します。このメトリックに基づいてこのパラメーターを調整できます。

全量データ読み取りフェーズでは、通常、各バッチで読み取る行数を減らす必要があります。 scan.incremental.snapshot.chunk.sizeパラメータの値を減らすことができます。

説明

VVR 11.5 以降でのみサポートされています。

include-binlog-meta.enable

メッセージに、GTID やバイナリログオフセットなどの元の MySQL バイナリログ情報を含めるかどうかを指定します。

いいえ

BOOLEAN

false

これは、既存の Canal 同期リンクを置き換えるなど、元のバイナリログ同期シナリオに適用できます。

説明

VVR 11.6 以降でのみサポートされています。

scan.binlog.tolerate.gtid-holes

このパラメータを有効にすると、GTID シーケンスのギャップを無視し、ジョブが不連続なイベントをバイパスして実行を継続できるようになります。

いいえ

BOOLEAN

false

このパラメータを有効にする前に、ジョブの開始オフセットが期限切れになっていないことを確認する必要があります。 ジョブが削除された、または期限切れの GTID オフセットから開始された場合、エンジンは欠落したログを警告なしにスキップするため、データ損失につながります。

説明

このパラメータは VVR 11.6 以降でのみサポートされています。

scan.emit.create-table-events.in-batch.enabled

ジョブの初期化フェーズ中にテーブルスキーマをバッチ送信するかどうかを指定します。

いいえ

BOOLEAN

false

これは実験的な機能です。 単一のジョブが多数のテーブルを同期する場合に、このオプションを有効にすることを推奨します。

説明

このパラメータは VVR 11.4 以降でのみサポートされています。

既存カタログの再利用

VVR 11.5 以降では、Flink CDC データインジェスト ジョブで、[Data Management] ページで作成した組み込み MySQL カタログを直接参照できます。これにより、接続プロパティを手動で記述する手間を削減できます。

source:
  type: mysql
  using.built-in-catalog: mysql_rds_catalog

現在、データインジェスト ジョブは、次の MySQL カタログ パラメーターを自動的に再利用できます:

  • hostname

  • port

  • username

  • password

  • catalog.table.metadata-columns

  • catalog.table.treat-tinyint1-as-boolean

これらの自動的に再利用されるパラメーターをオーバーライドする場合は、対応する YAML パラメーターを明示的に記述できます。明示的に記述したパラメーターの方が優先度が高くなります。

型マッピング

次の表に、データインジェストにおけるデータ型マッピングを示します。

MySQL CDC フィールドタイプ

CDC フィールドタイプ

TINYINT(n)

TINYINT

SMALLINT

SMALLINT

TINYINT UNSIGNED

TINYINT UNSIGNED ZEROFILL

YEAR

サポートされていません

INT

INT

MEDIUMINT

MEDIUMINT UNSIGNED

MEDIUMINT UNSIGNED ZEROFILL

SMALLINT UNSIGNED

SMALLINT UNSIGNED ZEROFILL

BIGINT

BIGINT

INT UNSIGNED

INT UNSIGNED ZEROFILL

BIGINT UNSIGNED

DECIMAL(20, 0)

BIGINT UNSIGNED ZEROFILL

SERIAL

FLOAT [UNSIGNED] [ZEROFILL]

FLOAT

DOUBLE [UNSIGNED] [ZEROFILL]

DOUBLE

DOUBLE PRECISION [UNSIGNED] [ZEROFILL]

REAL [UNSIGNED] [ZEROFILL]

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] and p <= 38

DECIMAL(p, s)

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] and p <= 38

FIXED(p, s) [UNSIGNED] [ZEROFILL] and p <= 38

BOOLEAN

BOOLEAN

BIT(1)

TINYINT(1)

DATE

DATE

TIME [(p)]

TIME [(p)]

DATETIME [(p)]

TIMESTAMP [(p)]

TIMESTAMP [(p)]

マッピングは、treat-timestamp-as-datetime-enabled パラメーターの値によって異なります:

true :TIMESTAMP[(p)]

false :TIMESTAMP_LTZ[(p)]

CHAR(n)

CHAR(n)

VARCHAR(n)

VARCHAR(n)

BIT(n)

BINARY(⌈(n + 7) / 8⌉)

BINARY(n)

BINARY(n)

VARBINARY(N)

VARBINARY(N)

NUMERIC(p, s) [UNSIGNED] [ZEROFILL] and 38 < p <= 65

STRING

説明

MySQL では、10進数データ型の精度は最大 65 ですが、Flink では精度は 38 に制限されています。 したがって、精度が 38 を超える10進数カラムを定義する場合は、精度が損なわれるのを避けるために文字列にマッピングする必要があります。

DECIMAL(p, s) [UNSIGNED] [ZEROFILL] and 38 < p <= 65

FIXED(p, s) [UNSIGNED] [ZEROFILL] and 38 < p <= 65

TINYTEXT

STRING

TEXT

MEDIUMTEXT

LONGTEXT

ENUM

JSON

STRING

説明

JSON データ型は、Flink では JSON 形式の文字列に変換されます。

GEOMETRY

STRING

説明

MySQL の空間データ型は、固定 JSON 形式の文字列に変換されます。 詳細については、MySQL の「空間データ型マッピング」をご参照ください。

POINT

LINESTRING

POLYGON

MULTIPOINT

MULTILINESTRING

MULTIPOLYGON

GEOMETRYCOLLECTION

TINYBLOB

BYTES

説明

MySQL の BLOB データ型では、長さが 2,147,483,647 (2**31-1) 以下の BLOB のみがサポートされます。

BLOB

MEDIUMBLOB

LONGBLOB

使用例

  • CDC ソーステーブル

    CREATE TEMPORARY TABLE mysqlcdc_source (
       order_id INT,
       order_date TIMESTAMP(0),
       customer_name STRING,
       price DECIMAL(10, 5),
       product_id INT,
       order_status BOOLEAN,
       PRIMARY KEY(order_id) NOT ENFORCED
    ) WITH (
      'connector' = 'mysql',
      'hostname' = '<yourHostname>',
      'port' = '3306',
      'username' = '<yourUsername>',
      'password' = '<yourPassword>',
      'database-name' = '<yourDatabaseName>',
      'table-name' = '<yourTableName>'
    );
    
    CREATE TEMPORARY TABLE blackhole_sink(
      order_id INT,
      customer_name STRING
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO blackhole_sink
    SELECT order_id, customer_name FROM mysqlcdc_source;
  • ディメンションテーブル

    CREATE TEMPORARY TABLE datagen_source(
      a INT,
      b BIGINT,
      c STRING,
      `proctime` AS PROCTIME()
    ) WITH (
      'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE mysql_dim (
      a INT,
      b VARCHAR,
      c VARCHAR
    ) WITH (
      'connector' = 'mysql',
      'hostname' = '<yourHostname>',
      'port' = '3306',
      'username' = '<yourUsername>',
      'password' = '<yourPassword>',
      'database-name' = '<yourDatabaseName>',
      'table-name' = '<yourTableName>'
    );
    
    CREATE TEMPORARY TABLE blackhole_sink(
      a INT,
      b STRING
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO blackhole_sink
    SELECT T.a, H.b
    FROM datagen_source AS T JOIN mysql_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H ON T.a = H.a;
  • シンクテーブル

    CREATE TEMPORARY TABLE datagen_source (
      `name` VARCHAR,
      `age` INT
    ) WITH (
      'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE mysql_sink (
      `name` VARCHAR,
      `age` INT
    ) WITH (
      'connector' = 'mysql',
      'hostname' = '<yourHostname>',
      'port' = '3306',
      'username' = '<yourUsername>',
      'password' = '<yourPassword>',
      'database-name' = '<yourDatabaseName>',
      'table-name' = '<yourTableName>'
    );
    
    INSERT INTO mysql_sink
    SELECT * FROM datagen_source;
  • データインジェスト用のデータソース (MySQL CDC)

    source:
      type: mysql
      name: MySQL Source
      hostname: ${mysql.hostname}
      port: ${mysql.port}
      username: ${mysql.username}
      password: ${mysql.password}
      tables: ${mysql.source.table}
      server-id: 7601-7604
    
    sink:
      type: values
      name: Values Sink
      print.enabled: true
      sink.print.logger: true

MySQL CDC ソーステーブルについて

  • 仕組み

    MySQL CDC ソーステーブルが起動すると、テーブル全体をスキャンし、プライマリキーに基づいてテーブルを複数のチャンクに分割し、現在のバイナリログオフセットを記録します。その後、ソーステーブルは増分スナップショットアルゴリズムを使用し、SELECT ステートメントで各チャンクのデータを読み取ります。ジョブは定期的にチェックポイントを実行して、完了したチャンクを記録します。フェールオーバーが発生した場合、ジョブは未完了のチャンクからデータの読み取りを継続します。すべてのチャンクの読み取りが完了すると、ジョブは、事前に記録したバイナリログオフセットから増分変更レコードの読み取りを開始します。Flink ジョブは引き続き定期的にチェックポイントを実行して、バイナリログオフセットを記録します。ジョブでフェールオーバーが発生した場合、最後に記録されたバイナリログオフセットから処理を再開し、exactly-once セマンティクスを実現します。

    増分スナップショットアルゴリズムの詳細については、「MySQL CDC Connector」をご参照ください。

  • メタデータ

    メタデータは、シャード化されたデータベースとテーブルのデータをマージして同期するシナリオで役立ちます。これは、マージ後に、ビジネス側で各データレコードのソースデータベースとテーブルを識別したいケースが多いためです。Metadata columns を使用すると、ソーステーブルのデータベース名とテーブル名の情報にアクセスできます。そのため、メタデータカラムを使用して複数のシャード化されたテーブルを 1 つの宛先テーブルに簡単にマージできます。

    MySQL CDC ソースはメタデータカラムの構文をサポートしています。メタデータカラムを使用して、次のメタデータにアクセスできます。

    メタデータキー

    メタデータタイプ

    説明

    database_name

    STRING NOT NULL

    行を含むデータベースの名前。

    table_name

    STRING NOT NULL

    行を含むテーブルの名前。

    op_ts

    TIMESTAMP_LTZ(3) NOT NULL

    データベース内で行が変更された時刻。レコードがバイナリログではなくテーブルの履歴データから取得された場合、この値は常に 0 です。

    説明

    このフィールドの精度は秒単位です。

    op_type

    STRING NOT NULL

    行の変更タイプ。

    • +I:INSERT メッセージ

    • -D:DELETE メッセージ

    • -U:UPDATE_BEFORE メッセージ

    • +U:UPDATE_AFTER メッセージ

    説明

    VVR 8.0.7 以降でのみサポートされています。

    query_log

    STRING NOT NULL

    この行の MySQL クエリログレコードを読み取ることができます。

    説明

    クエリログを記録するには、MySQL で binlog_rows_query_log_events パラメーターを有効にする必要があります。

    次のコード例は、MySQL インスタンス内の複数のシャード化されたデータベースにある複数の orders テーブルを、Hologres の holo_orders テーブルにマージして同期する方法を示しています。

    CREATE TEMPORARY TABLE mysql_orders (
      db_name STRING METADATA FROM 'database_name' VIRTUAL,  -- データベース名を読み取ります。
      table_name STRING METADATA  FROM 'table_name' VIRTUAL, -- テーブル名を読み取ります。
      operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL, -- 変更時間を読み取ります。
      op_type STRING METADATA FROM 'op_type' VIRTUAL, -- 変更タイプを読み取ります。
      order_id INT,
      order_date TIMESTAMP(0),
      customer_name STRING,
      price DECIMAL(10, 5),
      product_id INT,
      order_status BOOLEAN,
      PRIMARY KEY(order_id) NOT ENFORCED
    ) WITH (
      'connector' = 'mysql-cdc',
      'hostname' = 'localhost',
      'port' = '3306',
      'username' = 'flinkuser',
      'password' = 'flinkpw',
      'database-name' = 'mydb_.*', -- 複数のシャード化されたデータベースに一致する正規表現。
      'table-name' = 'orders_.*'   -- 複数のシャード化されたテーブルに一致する正規表現。
    );
    
    INSERT INTO holo_orders SELECT * FROM mysql_orders;

    上記のコードで、WITH 句で scan.read-changelog-as-append-only.enabled パラメーターを true に設定した場合、出力結果はダウンストリームテーブルのプライマリキー設定によって異なります:

    • ダウンストリームテーブルのプライマリキーが order_id の場合、出力結果にはアップストリームテーブル内の各プライマリキーに対する最後の変更のみが含まれます。あるプライマリキーに対する最後の変更が削除操作の場合、同じプライマリキーを持ち、op_type が -D のレコードをダウンストリームテーブルで確認できます。

    • ダウンストリームテーブルのプライマリキーが order_idoperation_ts、および op_type の場合、出力結果にはアップストリームテーブル内の各プライマリキーに対する完全な変更が含まれます。

  • 正規表現のサポート

    MySQL CDC ソーステーブルでは、テーブル名またはデータベース名に正規表現を使用し、複数のテーブルまたはデータベースに一致させることができます。次のコード例は、正規表現を使用して複数のテーブルを指定する方法を示しています。

    CREATE TABLE products (
      db_name STRING METADATA FROM 'database_name' VIRTUAL,
      table_name STRING METADATA  FROM 'table_name' VIRTUAL,
      operation_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,
      order_id INT,
      order_date TIMESTAMP(0),
      customer_name STRING,
      price DECIMAL(10, 5),
      product_id INT,
      order_status BOOLEAN,
      PRIMARY KEY(order_id) NOT ENFORCED
    ) WITH (
      'connector' = 'mysql-cdc',
      'hostname' = 'localhost',
      'port' = '3306',
      'username' = 'root',
      'password' = '123456',
      'database-name' = '(^(test).*|^(tpc).*|txc|.*[p$]|t{2})', -- 複数のデータベースに一致する正規表現。
      'table-name' = '(t[5-8]|tt)' -- 複数のテーブルに一致する正規表現。
    );

    例の正規表現の意味は次のとおりです:

    • ^(test).* は前方一致の例です。この式は、"test" で始まる "test1" や "test2" などのデータベース名に一致します。

    • .*[p$] は後方一致の例です。この式は、"p" で終わる "cdcp" や "edcp" などのデータベース名に一致します。

    • txc は完全一致です。"txc" と完全に一致するデータベース名に一致します。

    MySQL CDC が完全修飾テーブル名に一致させる場合、database-name.table-name パターンを使用してテーブルを一意に識別します。たとえば、パターン (^(test).*|^(tpc).*|txc|.*[p$]|t{2}).(t[5-8]|tt) は、データベース内の txc.tttest2.test5 などのテーブルに一致します。

    重要

    SQL ジョブの設定では、table-name および database-name パラメーターでカンマ (,) を使用して複数のテーブルまたはデータベースを指定することはできません。

    • 複数のテーブルに一致させる場合、または複数の正規表現を使用する場合は、縦棒 (|) で連結し、括弧で囲みます。たとえば、userproduct テーブルを読み取る場合は、table-name(user|product) に設定できます。

    • 正規表現にカンマが含まれる場合は、縦棒 (|) 演算子を使用して書き換える必要があります。たとえば、正規表現 mytable_\d{1,2} は、カンマを含まない同等の (mytable_\d{1}|mytable_\d{2}) に書き換える必要があります。

  • 同時実行制御

    MySQL コネクタは、全量データのマルチスレッド読み取りをサポートしており、データロード効率を向上させることができます。また、Realtime Compute for Apache Flink コンソールの Autopilot 自動チューニング機能と組み合わせることで、マルチスレッド読み取りの完了後に増分フェーズで自動的にスケールインし、コンピューティングリソースを節約できます。

    Realtime Compute for Apache Flink の開発コンソールでは、[Resource Configuration] ページで、基本モードまたはエキスパートモードでジョブの同時実行数を設定できます。違いは次のとおりです:

    • 基本モードで設定する同時実行数は、ジョブ全体のグローバルな同時実行数です。基础模式

    • エキスパートモードでは、必要に応じて特定の VERTEX の同時実行数を設定できます。vertex并发

    リソース構成の詳細については、「Configure deployment information for a job」をご参照ください。

    重要

    基本モードでもエキスパートモードでも、同時実行数を設定する場合、テーブルで宣言されるサーバー ID 範囲はジョブの同時実行数以上である必要があります。たとえば、サーバー ID 範囲が 5404-5412 の場合、一意のサーバー ID は 9 個あります。そのため、ジョブの同時実行数は最大 9 に設定できます。同じ MySQL インスタンスに対する異なるジョブでは、サーバー ID 範囲が重複してはいけません。つまり、各ジョブで異なるサーバー ID またはサーバー ID 範囲を明示的に設定する必要があります。

  • Autopilot による自動スケールイン

    全量データフェーズでは、大量の履歴データが蓄積されます。読み取り効率を向上させるため、通常は履歴データを並列に読み取ります。増分バイナリログフェーズでは、バイナリログデータ量が少なく、かつグローバル順序を確保する必要があるため、通常はシングルスレッド読み取りで十分です。自動チューニング機能を使用すると、全量フェーズと増分フェーズで異なるリソース要件を、パフォーマンスとリソースの観点でバランスを取ることができます。

    自動チューニングは、MySQL CDC ソースの各タスクのトラフィックを監視します。バイナリログフェーズに入ると、1 つのタスクのみがバイナリログ読み取りを担当し、他のタスクがアイドル状態になると、自動チューニングはソースの CU 数と同時実行数を自動的に削減します。自動チューニングを有効にするには、ジョブの O&M ページで自動チューニングモードを Active に設定します。

    説明

    同時実行数を削減するためのデフォルトの最小トリガー間隔は 24 時間です。自動チューニングのパラメーターと詳細については、「Configure automatic tuning」をご参照ください。

  • 起動モード

    scan.startup.mode 設定項目を使用して、MySQL CDC ソーステーブルの起動モードを指定します。オプションは次のとおりです:

    • initial (default):初回起動時にデータベーステーブルの全量読み取りを実行し、その後、増分モードに切り替えてバイナリログを読み取ります。

    • earliest-offset:スナップショットフェーズをスキップし、利用可能な最も古いバイナリログオフセットから読み取りを開始します。

    • latest-offset:スナップショットフェーズをスキップし、バイナリログの末尾から読み取りを開始します。このモードでは、ソーステーブルはジョブの開始後に発生したデータ変更のみを読み取ることができます。

    • specific-offset:スナップショットフェーズをスキップし、指定したバイナリログオフセットから読み取りを開始します。オフセットは、バイナリログのファイル名と位置、または GTID セットで指定できます。

    • timestamp:スナップショットフェーズをスキップし、指定したタイムスタンプからバイナリログイベントの読み取りを開始します。

    使用例:

    CREATE TABLE mysql_source (...) WITH (
        'connector' = 'mysql-cdc',
        'scan.startup.mode' = 'earliest-offset', -- 最も古いオフセットから開始します。
        'scan.startup.mode' = 'latest-offset', -- 最新のオフセットから開始します。
        'scan.startup.mode' = 'specific-offset', -- 特定のオフセットから開始します。
        'scan.startup.mode' = 'timestamp', -- 特定のタイムスタンプから開始します。
        'scan.startup.specific-offset.file' = 'mysql-bin.000003', -- specific-offset モードでバイナリログのファイル名を指定します。
        'scan.startup.specific-offset.pos' = '4', -- specific-offset モードでバイナリログの位置を指定します。
        'scan.startup.specific-offset.gtid-set' = '24DA167-0C0C-11E8-8442-00059A3C7B00:1-19', -- specific-offset モードで GTID セットを指定します。
        'scan.startup.timestamp-millis' = '1667232000000' -- timestamp モードで起動タイムスタンプを指定します。
        ...
    )
    重要
    • MySQL ソースは、チェックポイント中に現在のオフセットを INFO レベルでログに出力します。ログのプレフィックスは Binlog offset on checkpoint {checkpoint-id} です。このログは、特定のチェックポイントオフセットからジョブを開始する際に役立ちます。

    • 読み取るテーブルでスキーマ変更が発生した場合、earliest-offsetspecific-offset、または timestamp から開始するとエラーが発生する可能性があります。これは、Debezium リーダーは内部的に最新のテーブルスキーマを保存するため、スキーマが一致しない古いデータを正しくパースできないからです。

  • プライマリキーのない CDC ソーステーブルについて

    • プライマリキーのないテーブルを使用するには、scan.incremental.snapshot.chunk.key-column を設定する必要があり、非 NULL カラムのみを選択できます。

    • プライマリキーのない CDC ソーステーブルの処理セマンティクスは、scan.incremental.snapshot.chunk.key-column で指定したカラムの動作によって決まります:

      • 指定したカラムが更新されない場合、exactly-once セマンティクスを保証できます。

      • 指定したカラムが更新される場合、at-least-once セマンティクスのみを保証できます。ただし、ダウンストリーム側でプライマリキーを指定し、べき等操作を使用することで、データの正確性を確保できます。

  • Alibaba Cloud ApsaraDB RDS for MySQL のバックアップログの読み取り

    MySQL CDC ソーステーブルは、Alibaba Cloud ApsaraDB RDS for MySQL のバックアップログの読み取りをサポートしています。これは、全量データフェーズに時間がかかり、ローカルのバイナリログファイルが自動的に削除されてしまったものの、自動または手動でアップロードされたバックアップファイルは残っている、といったシナリオで役立ちます。

    使用例:

    CREATE TABLE mysql_source (...) WITH (
        'connector' = 'mysql-cdc',
        'rds.region-id' = 'cn-beijing',
        'rds.access-key-id' = 'xxxxxxxxx', 
        'rds.access-key-secret' = 'xxxxxxxxx', 
        'rds.db-instance-id' = 'rm-xxxxxxxxxxxxxxxxx', 
        'rds.main-db-id' = '12345678',
        'rds.download.timeout' = '60s'
        ...
    )
  • CDC ソースの再利用の有効化

    同じジョブ内で複数の MySQL CDC ソーステーブルを使用すると、複数のバイナリログクライアントが起動します。すべてのソーステーブルが同じインスタンス上にある場合、データベースの負荷が増加します。詳細については、「MySQL CDC FAQ」をご参照ください。

    解決策

    VVR 8.0.7 以降では、MySQL CDC ソースの再利用に対応しています。この機能により、マージ可能な MySQL CDC ソーステーブルがマージされます。マージは、データベース名、テーブル名、および server-id を除き、ソーステーブルの設定が同一である場合に行われます。エンジンは、同じジョブ内の MySQL CDC ソースを自動的にマージします。

    手順

    1. SQL ジョブで SET コマンドを使用します:

      SET 'table.optimizer.source-merge.enabled' = 'true';
      
      # (VVR 8.0.8 および 8.0.9 の場合) この項目も設定します:
      SET 'sql-gateway.exec-plan.enabled' = 'false';
      VVR 11.1 以降では、再利用はデフォルトで有効になっています。
    2. state なしでジョブを開始してください。 ソースの再利用設定を変更するとジョブトポロジーが変わるため、state なしでジョブを開始する必要があります。そうしない場合、ジョブが起動できないか、データ損失が発生する可能性があります。ソースがマージされると、トポロジーに MergetableSourceScan ノードが表示されます。

    重要
    • 再利用を有効にした後、オペレーターチェーンを無効にしないでください。pipeline.operator-chainingfalse に設定すると、データのシリアル化と逆シリアル化のオーバーヘッドが増加します。マージされるソースが多いほど、オーバーヘッドは大きくなります。

    • VVR 8.0.7 では、オペレーターチェーンを無効にするとシリアライズの問題が発生します。

バイナリログ読み取りの高速化

MySQL コネクタをソーステーブルまたはデータインジェストのデータソースとして使用する場合、インクリメンタルフェーズ中にバイナリログファイルをパースしてさまざまな変更メッセージを生成します。 バイナリログファイルは、すべてのテーブルの変更をバイナリ形式で記録します。 次の方法でバイナリログファイルのパースを高速化できます。

  • パースフィルター設定を有効にする

    • scan.only.deserialize.captured.tables.changelog.enabled 設定項目を使用して、指定されたテーブルの変更イベントのみをパースします。

  • Debezium パラメーターの最適化

    debezium.max.queue.size: 162580
    debezium.max.batch.size: 40960
    debezium.poll.interval.ms: 50
    • debezium.max.queue.size: ブロッキングキューが保持できるレコードの最大数。 Debezium がデータベースからイベントストリームを読み取ると、ダウンストリームに書き込む前にイベントをブロッキングキューに配置します。 デフォルト値は 8192 です。

    • debezium.max.batch.size: コネクタが各イテレーションで処理するイベントの最大数。 デフォルト値は 2048 です。

    • debezium.poll.interval.ms: コネクタが新しい変更イベントを要求する前に待機するミリ秒数。 デフォルト値は 1000 ミリ秒、つまり 1 秒です。

使用例:

CREATE TABLE mysql_source (...) WITH (
    'connector' = 'mysql-cdc',
    -- Debezium 設定
    'debezium.max.queue.size' = '162580',
    'debezium.max.batch.size' = '40960',
    'debezium.poll.interval.ms' = '50',
    -- パースフィルターを有効にする
    'scan.only.deserialize.captured.tables.changelog.enabled' = 'true',  -- 指定されたテーブルの変更イベントのみをパースします。
    ...
)
source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: ${mysql.source.table}
  server-id: 7601-7604
  # Debezium 設定
  debezium.max.queue.size: 162580
  debezium.max.batch.size: 40960
  debezium.poll.interval.ms: 50
  # パースフィルターを有効にする
  scan.only.deserialize.captured.tables.changelog.enabled: true

MySQL CDC Enterprise Edition のバイナリログ処理能力は 85 MB/s で、オープンソースコミュニティ版の約 2 倍です。 バイナリログファイルの生成速度が 85 MB/s (つまり、6 秒ごとに 512 MB のファイル 1 つ) を超えると、Flink ジョブのレイテンシーは増加し続けます。 バイナリログファイルの生成速度が遅くなると、処理レイテンシーは徐々に減少します。 バイナリログファイルに大きなトランザクションが含まれている場合、処理レイテンシーが一時的に増加する可能性があります。 そのトランザクションのログが読み取られた後、レイテンシーは減少します。

MySQL CDC DataStream API

重要

DataStream API を使用してデータの読み取りと書き込みを行う場合は、対応する DataStream コネクタを使用して Apache Flink に接続する必要があります。DataStream コネクタの設定方法の詳細については、「DataStream コネクタの使用方法」をご参照ください。

DataStream API プログラムを作成し、MySqlSource を使用できます。以下にコードと pom 依存関係の例を示します。

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
public class MySqlSourceExample {
  public static void main(String[] args) throws Exception {
    MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
        .hostname("yourHostname")
        .port(yourPort)
        .databaseList("yourDatabaseName") // キャプチャするデータベースを設定
        .tableList("yourDatabaseName.yourTableName") // キャプチャするテーブルを設定
        .username("yourUsername")
        .password("yourPassword")
        .deserializer(new JsonDebeziumDeserializationSchema()) // SourceRecord を JSON 形式の文字列に変換
        .build();
    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    // チェックポイントを有効にする
    env.enableCheckpointing(3000);
    env
      .fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
      // 4 つの並列ソースタスクを設定
      .setParallelism(4)
      .print().setParallelism(1); // メッセージの順序を保つため、シンクの並列度を 1 にします
    env.execute("Print MySQL Snapshot + Binlog");
  }
}
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-core</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-streaming-java</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-base</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-common</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-mysql</artifactId>
    <version>${vvr.version}</version>
</dependency>

MySqlSource を構築する際、コード内で次のパラメーターを指定する必要があります。

パラメーター

説明

hostname

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

port

MySQL データベースサービスのポート番号。

databaseList

MySQL データベースの名前。

説明

このパラメーターは、複数のデータベースからデータを読み取るための正規表現をサポートします。 .* を使用して、すべてのデータベースに一致させることができます。

username

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

password

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

deserializer

SourceRecord 型のレコードを指定された型にデシリアライズするデシリアライザー。有効な値:

  • RowDataDebeziumDeserializeSchema:SourceRecord を Apache Flink Table または SQL の内部データ構造 RowData に変換します。

  • JsonDebeziumDeserializationSchema:SourceRecord を JSON 形式の文字列に変換します。

pom 依存関係では、次のパラメーターを指定する必要があります。

${vvr.version}

Alibaba Cloud Realtime Compute for Apache Flink のエンジンバージョン、例:1.17-vvr-8.0.4-3

説明

ホットフィックスは、他のチャネルでは通知されずにリリースされることがあるため、Maven に表示されているバージョン番号を使用してください。

${flink.version}

Apache Flink のバージョン。例: 1.17.2

重要

ジョブ実行時の非互換性の問題を回避するため、Realtime Compute for Apache Flink のエンジンバージョンに対応する Apache Flink のバージョンを使用してください。バージョンマッピングの詳細については、「エンジンバージョン」をご参照ください。

よくある質問

CDC ソーステーブルの使用時に発生する可能性がある問題の詳細については、「CDC FAQ」をご参照ください。