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

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

最終更新日:Sep 18, 2026

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

前提条件

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 コネクタは、バイナリログトランザクション圧縮機能をサポートしていません。したがって、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 を設定することもできますが、これにより同期性能が犠牲になります。

注意事項

  • サーバー ID を設定してバイナリログ消費の競合を回避する。

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

データインジェスト

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

構文

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.file と scan.startup.specific-offset.pos の両方を設定するか、scan.startup.specific-offset.gtid-set のみを設定して特定の GTID セットから開始することでオフセットを指定できます。

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

重要

earliest-offset、specific-offset、timestamp の起動モードでは、起動時のテーブルスキーマが指定された開始オフセット時のスキーマと異なる場合、スキーマの不一致によりジョブはエラーを報告します。つまり、これら 3 つの起動モードを使用する場合、指定されたバイナリログ消費位置とジョブ起動時の間で、対応するテーブルのスキーマが変更されていないことを確認する必要があります。

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 の時間的な値」をご参照ください。

scan.startup.specific-offset.skip-events

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

いいえ

INTEGER

なし

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

scan.startup.specific-offset.skip-rows

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

いいえ

INTEGER

なし

このパラメーターを使用する場合、scan.startup.mode を specific-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.enabled を true に設定する必要があります。

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_ts、es_ts、query_log、file、pos があります。カンマを使用して複数のメタデータ列を区切ることができます。

説明

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

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

重要

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

説明

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

treat-tinyint1-as-boolean.enabled

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

いいえ

BOOLEAN

true

有効な値:

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

  • false:TINYINT(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 型はリテラルな時刻を格納し、タイムゾーンの影響を受けません。

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

include-comments.enabled

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

いいえ

BOOELEAN

false

有効な値:

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

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

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

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

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

サーバー ID を設定してバイナリログ消費の競合を回避

データインジェストジョブがバイナリログを読み取る際、ソースは server-id を使用して MySQL にレプリケーションクライアントとして登録します。複数のジョブまたは他のレプリケーションクライアントが同じサーバー ID を使用すると、バイナリログ消費の競合が発生し、ジョブが失敗します。サーバー ID を設定する際には、以下の点に注意してください:

  • デフォルトでは、server-id は 5400 から 6400 までのランダムな単一の値です。複数のジョブがデフォルト値を使用すると、競合が発生する可能性があります。ジョブごとに重複しない ID を明示的に設定することを推奨します。

  • ソースの並列度が 1 より大きい場合は、サーバー ID の範囲を設定する必要があります。範囲内の利用可能な ID の数は並列度以上でなければならず、各並列リーダーは異なる ID を使用します。

例:ソースの並列度が 4 の場合、4 つの ID を含む範囲を設定します。

source:
  type: mysql
  name: MySQL Source
  hostname: <hostname>
  port: 3306
  username: <username>
  password: <password>
  tables: app_db.\.*
  server-id: 5400-5403

sink:
  type: hologres

複数のジョブが同じ MySQL インスタンスから読み取る場合は、各ジョブに重複しない範囲を割り当ててください。たとえば、ジョブ A は 5400-5403 を使用し、ジョブ B は 5404-5407 を使用します。

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

データインジェストのデータソースとして 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 秒です。

使用例:

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 configuration
  debezium.max.queue.size: 162580
  debezium.max.batch.size: 40960
  debezium.poll.interval.ms: 50
  # Enable parsing filter
  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 使用率メトリックと、JVM ページの TM GC 時間メトリックを確認して、CPU またはメモリリソースが不足していないか判断します。ジョブリソースを増やすことで、読み取り性能を最適化できます。

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

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

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

source:
  type: mysql
  hostname: <yourHostname>
  port: 3306
  username: <yourUsername>
  password: <yourPassword>
  tables: <yourTables>
  # RDS connection parameters for reading archived binlogs from OSS
  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 # Database instance ID.
  rds.main-db-id: 12345678 # Primary database ID.
  rds.endpoint: rds.aliyuncs.com

データベースとスキーマ同期のためのデータインジェスト

データ同期ロジックのみを含むジョブの場合、データインジェストジョブとして実行することを推奨します。データインジェストジョブは、データ統合シナリオ向けに深く最適化されています。使用方法については、「Flink CDC データインジェスト入門」および「Flink CDC データインジェストジョブの開発 (パブリックプレビュー)」をご参照ください。

以下のコードは、MySQL から Hologres へ app_db データベース全体を同期するデータインジェストジョブの例です。これには、上流データベースからの後続のスキーマ変更も含まれます:

source:
  type: mysql
  hostname: <hostname>
  port: 3306
  username: ${secret_values.mysqlusername}
  password: ${secret_values.mysqlpassword}
  tables: app_db.\.*
  server-id: 5400-5404

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <endpoint>
  dbname: <database-name>
  username: ${secret_values.holousername}
  password: ${secret_values.holopassword}

pipeline:
  name: Sync MySQL Database to Hologres

データインジェストにおける新規テーブルの検出

MySQL データインジェストコネクタは、2 つの異なるシナリオで新規テーブルの検出をサポートする設定オプションを提供します。

パラメーター

説明

注意

scan.newly-added-table.enabled

ジョブがチェックポイントから再起動する際、このオプションは前回の起動時に検出されなかったテーブルを同期します。これらの新しいテーブルからスナップショットデータと増分データの両方を読み取ります。

このオプションは、scan.startup.mode が initial に設定されている場合にのみサポートされます。他の起動モードでは効果がありません。

scan.binlog.newly-added-table.enabled

増分フェーズ中に、このオプションは新しく検出されたテーブルからデータを自動的に同期します。

  • ジョブを初めて開始する際にこのオプションを有効にすることを推奨します。ジョブは自動的に CREATE TABLE DDL 文を解析し、データを下流に同期します。データベーステーブルが既に作成された後にこのオプションを有効にしてジョブを再起動すると、データが不完全になる可能性があります。

  • initial 起動モードでは、スナップショットフェーズが完了するまで DDL 操作は下流に同期されません。スナップショットフェーズ中に作成されたテーブルは、scan.binlog.newly-added-table.enabled が有効になっていても自動的に同期できません。

重要
  • 完全読み取りフェーズ中に、セーブポイントを保存し、その後ソーステーブルを追加または削除してからセーブポイントから再起動することはサポートされていません。そうすると、ジョブがデータを正しく読み取れなくなります。

  • scan.newly-added-table.enabled と scan.binlog.newly-added-table.enabled を同時に有効にしないでください。両方を有効にすると、データが重複する可能性があります。