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

Realtime Compute for Apache Flink:MySQL SQL コネクタ

最終更新日:Sep 19, 2026

このトピックでは、SQL ジョブで MySQL コネクタを使用する方法について説明します。

背景情報

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

重要

MySQL コネクタを使用して OceanBase からデータを読み取る場合は、バイナリロギング (binlog) が有効になっており、正しく設定されていることを確認してください。詳細については、「Binlog 関連の操作」をご参照ください。この機能はパブリックプレビュー段階です。ご利用の際はご注意ください。

MySQL コネクタは以下をサポートします。

カテゴリ

詳細

サポートされているタイプ

ソーステーブル、ディメンションテーブル、結果テーブル、データインジェストデータソース

ランタイムモード

ストリーミングモードのみがサポートされています。

データフォーマット

該当なし

特定の監視メトリクス

監視メトリクス

  • ソーステーブル

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

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

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

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

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

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

説明

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

API タイプ

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

結果テーブルでのデータの更新または削除のサポート

はい

特徴

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

  • 完全データと増分データの両方の読み取りをサポートする統一されたバッチおよびストリーム処理により、2 つの別々のプロセスを維持する必要がなくなります。

  • 水平方向のパフォーマンススケーリングのための完全データの同時読み取り。

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

  • 完全データ読み取りフェーズ中の再開可能なデータ転送による安定性の向上。

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

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

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

前提条件

MySQL CDC ソーステーブルを使用する前に、「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 で導入され、デフォルトで無効になっています。

  • SELECT、SHOW DATABASES、REPLICATION SLAVE、および REPLICATION CLIENT 権限を持つ MySQL ユーザーが作成されていること。

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

  • IP アドレスホワイトリストを設定します。詳細については、「ApsaraDB RDS for MySQL インスタンスの IP アドレスホワイトリストの設定」をご参照ください。

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 で導入され、デフォルトで無効になっています。

  • SELECT、SHOW DATABASES、REPLICATION SLAVE、および REPLICATION CLIENT 権限を持つ MySQL ユーザーを作成していること。

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

  • IP アドレスホワイトリストを設定します。詳細については、「PolarDB for MySQL クラスターの IP アドレスホワイトリストの設定」をご参照ください。

セルフマネージド 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 データベースとテーブルを作成します。詳細については、「セルフマネージド MySQL インスタンスのデータベースとアカウントの作成」をご参照ください。権限不足による操作の失敗を防ぐため、特権アカウントを使用して MySQL データベースを作成してください。

  • IP アドレスホワイトリストを設定します。詳細については、「セルフマネージド MySQL インスタンスの IP アドレスホワイトリストの設定」をご参照ください。

制限事項

一般的な制限事項

  • MySQL CDC ソーステーブルは Watermark 定義をサポートしていません。

  • 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 以前のマルチマスタークラスターアーキテクチャのクラスターからのデータ読み取りをサポートしていません。詳細については、「マルチマスタークラスターとは」をご参照ください。これらのクラスターによって生成されるバイナリログには、重複したテーブル 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 を設定することもできますが、これによりレプリケーション性能が犠牲になります。

注意事項

  • ソーステーブル

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

  • シンクテーブル

    • 自動インクリメント主キー:DDL で自動インクリメント主キーを宣言しないでください。MySQL はデータ書き込み時に自動的にそれらを埋めます。

    • 少なくとも 1 つの非主キーフィールドを宣言する必要があります。そうしないと、エラーが報告されます。

    • DDL の `NOT ENFORCED` 制約は、Flink が主キー検証を強制しないことを示します。主キーの正確性と整合性を保証するのはユーザーの責任です。詳細については、「有効性チェック」をご参照ください。

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

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

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

SQL

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

構文

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

    なし

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

    説明

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

    username

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

    はい

    STRING

    なし

    なし。

    password

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

    はい

    STRING

    なし

    なし。

    database-name

    MySQL データベースの名前。

    はい

    STRING

    なし

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

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

    table-name

    MySQL テーブルの名前。

    はい

    STRING

    なし

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

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

    説明

    MySQL CDC ソーステーブルが正規表現を使用してテーブル名を照合する場合、指定した database-name と table-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 範囲を設定します。詳細については、「サーバー ID の使用」をご参照ください。

    scan.incremental.snapshot.enabled

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

    いいえ

    BOOLEAN

    true

    増分スナップショットはデフォルトで有効になっています。増分スナップショットは、完全なデータスナップショットを読み取るための新しいメカニズムです。古いスナップショット読み取り方法と比較して、増分スナップショットには次のような多くの利点があります:

    • ソースは完全なデータを並行して読み取ることができます。

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

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

    ソースが同時読み取りをサポートするようにしたい場合、各同時リーダーには一意のサーバー ID が必要です。したがって、server-id は 5400-6400 のような範囲でなければならず、範囲のサイズは同時実行数以上でなければなりません。

    説明

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

    scan.incremental.snapshot.chunk.size

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

    いいえ

    INTEGER

    8096

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

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

    scan.snapshot.fetch.size

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

    いいえ

    INTEGER

    1024

    なし。

    scan.startup.mode

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

    いいえ

    STRING

    initial

    有効な値:

    • initial (デフォルト):初回起動時またはステートレス起動時に、コネクタは完全な既存データをスキャンし、その後最新のバイナリログデータを読み取ります。

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

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

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

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

    重要

    earliest-offset、specific-offset、または timestamp 起動モードを使用する場合、指定されたバイナリログ消費位置とジョブ起動時間の間に対応するテーブルのスキーマが変更されないようにしてください。これにより、スキーマの不一致によるエラーを防ぎます。

    scan.startup.specific-offset.file

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

    いいえ

    STRING

    なし

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

    scan.startup.specific-offset.pos

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

    いいえ

    INTEGER

    なし

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

    scan.startup.specific-offset.gtid-set

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

    いいえ

    STRING

    なし

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

    scan.startup.timestamp-millis

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

    いいえ

    LONG

    なし

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

    重要

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

    server-time-zone

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

    いいえ

    STRING

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

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

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

    テーブルの行数がこの値より大きい場合、バッチ読み取りモードが使用されます。

    いいえ

    INTEGER

    1000

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

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

    • バッチ読み取り:データを複数のバッチで読み取ります。バッチごとに特定の行数で、すべてのデータが読み取られるまで続きます。この方法は、大きなテーブルを読み取る際の 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 設定プロパティ」をご参照ください。

    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

    なし

    • 主キーのないテーブルでは必須です。選択された列は非ヌル型 (NOT NULL) である必要があります。

    • 主キーのあるテーブルではオプションです。主キーから 1 つの列のみを選択できます。

    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

    なし

    説明

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

    rds.download.timeout

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

    いいえ

    DURATION

    60s

    なし。

    rds.endpoint

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

    いいえ

    STRING

    なし

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

    • 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 (デフォルト):スナップショット読み取りフェーズ中にバックフィルをスキップしません。

    バックフィルは単一チャンクのスナップショットクエリ中にのみ適用され、完全読み取りフェーズ全体をカバーするものではありません。バックフィルがスキップされると、各チャンクのスナップショットクエリはその時点での最新のテーブルデータを読み取ります。チャンクが読み取られた後に発生した更新は、完全読み取りフェーズ中にはマージされず、増分フェーズに入ってから Binlog から読み取られます。たとえば、chunk5 のスナップショット中に発生した chunk5 への更新は、chunk5 のスナップショットに直接反映されます。リーダーが chunk80 に進んだ後に chunk5 が更新された場合、その更新は後で増分フェーズ中に Binlog から適用されます。

    重要

    有効にすると、チャンクのスキャン中またはスキャン後に発生した変更は、増分フェーズで Binlog から配信され、重複する可能性があります。at-least-once セマンティクスのみが保証されます。下流の sink が主キーによるべき等書き込みをサポートしている場合にのみ有効にしてください。

    説明

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

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

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

    いいえ

    BOOELEAN

    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 以降でのみサポートされています。

    scan.filter-unchanged-row.enabled

    SQL で使用されていないフィールドの変更イベントをフィルタリングするかどうかを指定します。

    いいえ

    BOOLEAN

    false

    このパラメーターを true に設定すると、CDC ソーステーブルは UPDATE イベントを受信したときに、SQL ステートメントで参照されている列のみを比較します。参照されている列の値が変更されていない場合、イベントは下流に送信されません。このパラメーターは、複数ストリームの JOIN シナリオで有効です。上流テーブルが SQL ステートメントで参照されていないフィールドを更新した場合、不要な DELETE および INSERT イベントは生成されません。これにより、下流の演算子でのデータ増幅と結果のジッターが減少します。

    説明

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

    scan.projection-pushdown.enabled

    プロジェクションプッシュダウンを有効にするかどうかを指定します。

    いいえ

    BOOLEAN

    false

    このパラメーターを true に設定すると、CDC ソーステーブルは SQL ステートメントで参照されている列のみを解析して送信します。参照されていない列はシリアル化されず、ネットワーク経由で送信されません。これにより、I/O と計算のオーバーヘッドが削減されます。このパラメーターは scan.filter-unchanged-row.enabled と一緒に使用すると、より効果的です。プロジェクションプッシュダウンがまずデータ量を削減し、その後、変更されていない行がフィルタリングされて、下流の演算子の処理負荷がさらに軽減されます。

    説明

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

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

    パラメーター

    説明

    必須

    データの型

    デフォルト値

    注釈

    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

    キャッシュの生存時間 (TTL)。

    いいえ

    DURATION

    10 s

    lookup.cache.ttl の設定は lookup.cache.strategy に依存します:

    • lookup.cache.strategy が None に設定されている場合、lookup.cache.ttl を設定する必要はありません。これはキャッシュがタイムアウトしないことを意味します。

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

    • lookup.cache.strategy が ALL に設定されている場合、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 に設定する必要があります。そうしないと、ジョブが異常に実行される可能性があります。

    lookup.cache.cache-empty

    空のクエリ結果をキャッシュするかどうかを指定します。

    いいえ

    BOOLEAN

    true

    キャッシュポリシーが LRU の場合にのみ有効です。有効な値:

    • true (デフォルト):クエリがデータを返さない場合、空の結果がキャッシュされ、後続のクエリに対して直接返されます。

    • false:クエリがデータを返さない場合、結果はキャッシュされず、次のクエリで再度データが検索されます。

    説明

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

  • 結果テーブルのみ

    パラメーター

    説明

    必須

    データの型

    デフォルト値

    注釈

    url

    MySQL JDBC URL。

    いいえ

    STRING

    なし

    URL のフォーマットは次のとおりです:jdbc:mysql://<endpoint>:<port>/<database_name>。

    sink.max-retries

    データ書き込みに失敗した後の最大再試行回数。

    いいえ

    INTEGER

    3

    なし。

    sink.buffer-flush.batch-size

    単一のバッチ書き込みの行数。

    いいえ

    INTEGER

    4096

    なし。

    sink.buffer-flush.max-rows

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

    いいえ

    INTEGER

    10000

    このパラメーターは、主キーが指定された後にのみ有効になります。

    sink.buffer-flush.interval

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

    いいえ

    DURATION

    1s

    なし。

    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-delete-mode

    DELETE 操作が無視された後の delete タイプのレコードの処理戦略。

    いいえ

    STRING

    ALL

    有効な値:

    • ALL:-D と -U の両方のレコードを無視します。

    • REAL_DELETE:-D レコードのみを無視します。

    • UPDATE_BEFORE:-U レコードのみを無視します。

    説明
    • このオプションは Realtime Compute エンジン VVR 11.8 以降でのみサポートされています。

    • sink.ignore-delete=true の場合にのみ有効です。単独で設定するとエラーになります。

    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 以降でのみサポートされています。

    sink.force-batch-on-non-primary-table

    主キーのない MySQL 結果テーブルにデータが書き込まれるときに、レコードを強制的にバッファリングし、JDBC バッチを使用してバッチで挿入するかどうかを指定します。

    いいえ

    BOOLEAN

    false

    このパラメーターは、主キーのないテーブルでのみ有効です。このパラメーターを true に設定すると、DELETE および UPDATE_BEFORE イベントはサイレントに破棄され、INSERT および UPDATE_AFTER イベントは INSERT ステートメントとしてバッチで書き込まれます。データが追加専用であるか、前述のセマンティクスが許容できる場合にのみ、このパラメーターを有効にすることを推奨します。

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

    • データが追加専用であるか、前述のセマンティクスが許容できる場合にのみ、このパラメーターを有効にしてください。そうしないと、データが不整合になる可能性があります。

型マッピング

  • 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

使用例

  • 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;
  • データインジェストデータソース

    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 ジョブは定期的にチェックポイントを実行し続け、バイナリログオフセットを記録します。ジョブがフェールオーバーした場合、最後に記録されたバイナリログオフセットから処理を再開し、1 回限りのセマンティクスを実現します。

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

  • メタデータ

    メタデータは、シャーディングされたデータベースとテーブルのデータがマージおよび同期されるシナリオで役立ちます。これは、マージ後、ビジネスが各データレコードのソースデータベースとテーブルを区別したいことが多いためです。メタデータ列は、ソーステーブルのデータベースとテーブル名情報にアクセスできます。したがって、メタデータ列を使用して、複数のシャーディングされたテーブルを単一の宛先テーブルに簡単にマージできます。

    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_id`、`operation_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).*` はプレフィックスマッチングの例です。この式は、「test1」や「test2」など、「test」で始まるデータベース名に一致します。

    • `.*[p$]` はサフィックスマッチングの例です。この式は、「cdcp」や「edcp」など、「p」で終わるデータベース名に一致します。

    • `txc` は特定のマッチです。これは、正確に「txc」であるデータベース名に一致します。

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

    重要

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

    • 複数のテーブルを照合したり、複数の正規表現を使用したりするには、それらを縦棒 (|) で接続し、括弧で囲みます。たとえば、`user` テーブルと `product` テーブルを読み取るには、`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 の開発コンソールでは、リソース設定ページで基本モードまたはエキスパートモードでジョブの同時実行数を設定できます。

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

      たとえば、基本モードで parallelism が 8 に設定されている場合、SQL WITH 句の server-id は連続した範囲 (たとえば '404-412') として設定する必要があります。

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

    リソース設定の詳細については、「ジョブのデプロイメント情報の設定」をご参照ください。

    重要

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

  • Autopilot 自動スケールイン

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

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

    説明

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

  • 起動モード

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

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

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

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

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

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

    ステートレス起動は状態を再利用しません。ソースコネクタはこれを初回起動として扱うため、scan.startup.mode が再度有効になります。デプロイメントの起動モードの詳細については、「デプロイメントの開始」をご参照ください。

    使用例:

    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-offset`、`specific-offset`、または `timestamp` から開始するとエラーが発生する可能性があります。これは、Debezium リーダーが内部的に最新のテーブルスキーマを保存しており、スキーマが一致しない古いデータを正しく解析できないためです。

  • 主キーのない CDC ソーステーブルについて

    • 主キーのないテーブルを使用するには、`scan.incremental.snapshot.chunk.key-column` を設定する必要があり、非ヌル列のみを選択できます。

    • 主キーのない CDC ソーステーブルの処理セマンティクスは、`scan.incremental.snapshot.chunk.key-column` で指定された列の動作によって決まります:

      • 指定された列が更新されない場合、1 回限りのセマンティクスが保証されます。

      • 指定された列が更新される場合、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 よくある質問」をご参照ください。

    ソリューション

    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. 状態なしでジョブを開始します。ソースの再利用設定を変更するとジョブのトポロジーが変わるため、状態なしでジョブを開始する必要があります。そうしないと、ジョブの起動に失敗したり、データが失われたりする可能性があります。ソースがマージされた場合、トポロジーに MergetableSourceScan ノードが表示されます。

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

    • VVR 8.0.7 では、演算子チェーンを無効にするとシリアル化の問題が発生します。

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

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

  • 並列解析と解析フィルターを有効にする (この機能には、Ververica Runtime (VVR) 8.0.7 以降を搭載した Realtime Compute for Apache Flink が必要です。MySQL CDC コネクタのコミュニティ版では利用できません。)

    • scan.only.deserialize.captured.tables.changelog.enabled オプションを有効にして、指定されたテーブルの変更イベントのみを解析します。

    • scan.parallel-deserialize-changelog.enabled オプションを有効にして、複数のスレッドを使用してバイナリログファイルを解析し、イベントを順序通りにコンシューマーキューに配信します。このオプションを有効にする場合、通常は TaskManager CPU も増やす必要があります。

  • 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 を超える (つまり、512 MB のファイルが 6 秒ごとに 1 つ) と、Flink ジョブの遅延は増加し続けます。バイナリログファイルの生成速度が遅くなると、処理遅延は徐々に減少します。バイナリログファイルに大規模トランザクションが含まれている場合、処理遅延が一時的に増加することがあります。そのトランザクションのログが読み取られた後、遅延は減少します。

データ遅延の診断とジョブスループットの最適化

増分フェーズでデータ遅延が発生した場合は、次の手順で問題を分析してください:

  1. 概要ページの currentFetchEventTimeLag および currentEmitEventTimeLag メトリックを確認します。currentFetchEventTimeLag メトリックは、バイナリログからのデータ読み取りの遅延を表します。currentEmitEventTimeLag メトリックは、ジョブに関連するテーブルのバイナリログからのデータ読み取りの遅延を表します。

    シナリオ

    説明

    currentFetchEventTimeLag は低いが、currentEmitEventTimeLag は高く、ほとんど更新されない。

    低い currentFetchEventTimeLag は、データベースからのバイナリログのプルが効率的であることを示します。しかし、バイナリログにはジョブが読み取る必要のあるテーブルのデータがほとんど含まれていません。そのため、currentEmitEventTimeLag はほとんど更新されません。これは予期された動作です。

    currentFetchEventTimeLag と currentEmitEventTimeLag の両方が高い。

    これは、ソーステーブルの読み取りパフォーマンスが低いことを示します。このセクションの以降の手順で最適化を進めることができます。

  2. バックプレッシャーは、ソースが下流の演算子にデータを送信するレートを低下させる可能性があります。sourceIdleTime が定期的に増加し、currentFetchEventTimeLag と currentEmitEventTimeLag の両方が継続的に増加することが観察される場合があります。これを解決するには、バックプレッシャーが発生しているノードの並列度を上げてください。

  3. CPU ページの TM CPU Usage メトリックと JVM ページの TM GC Time メトリックを確認して、CPU またはメモリリソースが不足していないか判断します。ジョブリソースを増やすことで読み取りパフォーマンスを最適化できます。また、ミニバッチパラメーターを有効にしてスループットを向上させることもできます。詳細については、「高性能 Flink SQL 最適化技術」をご参照ください。

  4. ジョブに大きな状態を持つ SinkUpsertMaterializer 演算子が存在する場合、読み取りパフォーマンスに影響を与える可能性があります。ジョブの並列度を上げるか、SinkUpsertMaterializer 演算子を避けることを検討してください。詳細については、「SinkUpsertMaterializer の使用を避ける」をご参照ください。既存のジョブから SinkUpsertMaterializer 演算子を削除するには、ステートレス再起動が必要です。これは、ジョブトポロジーが変更されるため、既存の状態から開始するとジョブの失敗やデータ損失を引き起こす可能性があるためです。

サーバー ID を設定して binlog 競合を回避する

データベースからデータを同期する各クライアントには、サーバー ID と呼ばれる一意の ID があります。異なるジョブが同じサーバー ID を使用すると、競合が発生し、ジョブの失敗を引き起こす可能性があります。各 MySQL CDC データソースに異なるサーバー ID を割り当てることを推奨します。

  • サーバー ID の設定方法

    Flink テーブルの DDL ステートメントまたは SQL ヒントを使用してサーバー ID を指定できます。

    テーブル DDL の WITH 句で指定するのではなく、SQL ヒントを使用してサーバー ID を設定することを推奨します。詳細については、「SQL ヒント」をご参照ください。

  • 異なるシナリオでのサーバー ID 設定

    • 増分スナップショットが無効または並列度が 1 の場合

      増分スナップショットが無効または並列度が 1 の場合、単一のサーバー ID を指定できます。

      SELECT * FROM source_table /*+ OPTIONS('server-id'='123456') */ ;
    • 増分スナップショットが有効で並列度が 1 より大きい場合

      増分スナップショットが有効で並列度が 1 より大きい場合、サーバー ID の範囲を指定する必要があります。範囲内の利用可能なサーバー ID の数は、少なくとも並列度と同じでなければなりません。たとえば、並列度が 3 の場合、次の設定を使用できます:

      SELECT * FROM source_table /*+ OPTIONS('server-id'='123456-123458') */ ;
    • 複数の MySQL CDC ソーステーブルを含む Flink SQL ジョブ

      Flink SQL ジョブに複数の MySQL CDC ソーステーブルが含まれ、ソース再利用が無効になっている場合、各 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;

チャンクパラメーターを設定してメモリ使用量を最適化する

MySQL CDC ソーステーブルが起動すると、全表スキャンを実行し、主キーに基づいてテーブルを複数のチャンクに分割し、現在のバイナリログ位置を記録します。その後、ジョブは増分スナップショットアルゴリズムを使用して、SELECT ステートメントで各チャンクから順次データを読み取ります。ジョブは定期的にチェックポイントを実行して、完了したチャンクを記録します。フェールオーバーが発生した場合、最初の未完了チャンクから読み取りを再開します。すべてのチャンクが読み取られた後、ジョブは以前に記録されたバイナリログ位置から増分変更の読み取りに切り替わります。Flink ジョブは定期的にチェックポイントを実行してバイナリログ位置を保存します。フェールオーバーが発生した場合、ジョブは最後に保存された位置から処理を再開し、1 回限りのセマンティクスを実現します。

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

単一列の主キーを持つテーブルの場合、チャンクはデフォルトでそのキーに基づいて分割されます。複合主キーを持つテーブルの場合、デフォルトで主キーの最初の列が分割に使用されます。Ververica Runtime (VVR) 6.0.7 以降では、主キーのないソーステーブルの読み取りがサポートされています。scan.incremental.snapshot.chunk.key-column パラメーターを設定して、分割用の非ヌル列を指定する必要があります。

チャンクパラメーターの最適化

チャンクデータとメタデータはメモリに保存されるため、メモリ不足 (OOM) エラーが発生することがあります。OOM エラーが発生したコンポーネントに応じてパラメーターを調整できます:

  • JobManager

    JobManager はすべてのチャンクのメタデータを保存します。チャンク数が多すぎると OOM エラーが発生する可能性があります。これを解決するには、scan.incremental.snapshot.chunk.size の値を増やしてチャンク数を減らします。または、ランタイム設定で jobmanager.memory.heap.size を設定して JobManager のヒープメモリを増やすこともできます。詳細については、「Flink パラメーター設定」をご参照ください。

  • TaskManager

    • TaskManager は各チャンクのデータを読み取ります。チャンクに含まれる行が多すぎると、OOM エラーが発生する可能性があります。これを解決するには、scan.incremental.snapshot.chunk.size の値を減らしてチャンクあたりの行数を減らします。また、ランタイム設定で TaskManager Memory の値を増やして TaskManager のヒープメモリを増やすこともできます。

    • VVR 8.0.8 以前では、最後のチャンクに大量のデータが含まれる可能性があり、TaskManager で OOM エラーが発生することがあります。この問題を回避するために、VVR 8.0.9 以降へのアップグレードを推奨します。

    • 複合主キーを持つ MySQL CDC ソーステーブルの場合、チャンクはデフォルトでキーの最初の列に基づいて分割されます。データが大幅に偏っており、多くの行がその列で同じ値を共有している場合、その値のチャンクが非常に大きくなり、TaskManager で OOM エラーが発生する可能性があります。scan.incremental.snapshot.chunk.key-column を設定して、主キーから別の列を分割用に指定できます。

スナップショットフェーズでの読み取りの高速化

スナップショットフェーズ中、MySQL ソーステーブルは JDBC 接続を介してスナップショットデータを読み取ります。このフェーズでの読み取りを高速化するには、次の方法を使用します。

  1. ソースの並列度を上げて、スナップショットフェーズ中の読み取りを高速化します。

  2. scan.incremental.snapshot.chunk.size の値を増やして、単一のチャンクでより多くのデータを取得します。

  3. 下流の結果テーブルに主キーがあり、べき等書き込みをサポートしている場合は、scan.incremental.snapshot.backfill.skip を有効にして、バックフィル部分のバイナリログの読み取りをスキップできます。これにより、スナップショットフェーズ中の処理が高速化されます。

ソース再利用を有効にして binlog 接続を削減

ジョブに複数の MySQL ソーステーブルが含まれる場合、ソース再利用を有効にして、単一のバイナリログ接続を共有することでデータベースの負荷を軽減できます。この機能は Realtime Compute for Apache Flink でのみ利用可能であり、MySQL CDC コネクタのコミュニティ版ではサポートされていません。

SQL ジョブで SET コマンドを使用してソース再利用機能を有効にします:

SET 'table.optimizer.source-merge.enabled' = 'true';

ソース再利用は新しいジョブに対してのみ有効にすることを推奨します。既存のジョブでソース再利用を有効にする場合は、ステートレス再起動を実行する必要があります。これは、ソース再利用がジョブトポロジーを変更するため、既存の状態から開始するとジョブの失敗やデータ損失を引き起こす可能性があるためです。

ソース再利用を有効にすると、同じ設定パラメーターを持つ MySQL ソーステーブルがマージされます。ジョブ内のすべてのソーステーブルが同じ設定を共有している場合、バイナリログ接続の数は次のように計算されます:

  • スナップショットフェーズ中、バイナリログ接続の数はソースの並列度と同じです。

  • 増分フェーズ中、バイナリログ接続の数は 1 です。

重要
  • VVR 8.0.8 および 8.0.9 では、CDC ソースの再利用を有効にする際に SET 'sql-gateway.exec-plan.enabled' = 'false'; も設定する必要があります。

  • CDC ソースの再利用を有効にした後、pipeline.operator-chaining ジョブオプションを false に設定しないでください。演算子チェーンを壊すと、ソースから下流の演算子に送信されるデータのシリアル化と逆シリアル化のオーバーヘッドが増加します。マージされるソースが多いほど、オーバーヘッドは大きくなります。

  • Ververica Runtime (VVR) 8.0.7 では、pipeline.operator-chaining を false に設定するとシリアル化の問題が発生します。

OSS からのアーカイブログの読み取り

ApsaraDB RDS for MySQL インスタンスをデータソースとして使用する場合、OSS に保存されているログバックアップを読み取ることができます。指定されたタイムスタンプまたはバイナリログ位置に対応するファイルが OSS に保存されている場合、Flink は自動的に OSS からクラスターにログファイルをプルします。ファイルがデータベースにローカルに保存されている場合、Flink は自動的にデータベース接続を介した読み取りに切り替わります。この機能は Realtime Compute for Apache Flink でのみ利用可能であり、MySQL CDC コネクタのコミュニティ版ではサポートされていません。

OSS ログバックアップからの読み取りを有効にするには、ApsaraDB RDS for MySQL の接続パラメーターを設定する必要があります。例:

CREATE TABLE mysql_source (...) WITH (
    'connector' = 'mysql-cdc',
    'rds.region-id' = 'cn-beijing',
    'rds.access-key-id' = 'your_access_key_id',
    'rds.access-key-secret' = 'your_access_key_secret',
    'rds.db-instance-id' = 'rm-xxxxxxxx',  // データベースインスタンスの ID
    'rds.main-db-id' = '12345678', // プライマリデータベースの ID
    'rds.endpoint' = 'rds.aliyuncs.com'
    ...
)

よくある質問

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