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 を設定することもできますが、これにより同期性能が犠牲になります。
注意事項
完全データ読み取りフェーズ中に、セーブポイントを保存し、ソーステーブルに新しいテーブルを追加または削除してから、セーブポイントからジョブを再起動することはできません。これにより、ジョブがデータの読み取りに失敗します。
データインジェスト
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 | なし |
説明
|
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 | 有効な値:
重要 earliest-offset、specific-offset、timestamp の起動モードでは、起動時のテーブルスキーマが指定された開始オフセット時のスキーマと異なる場合、スキーマの不一致によりジョブはエラーを報告します。つまり、これら 3 つの起動モードを使用する場合、指定されたバイナリログ消費位置とジョブ起動時の間で、対応するテーブルのスキーマが変更されていないことを確認する必要があります。 |
scan.startup.specific-offset.file | specific-offset 起動モードを使用する場合の開始オフセットのバイナリログファイル名。 | いいえ | STRING | なし | このパラメーターを使用する場合、scan.startup.mode を specific-offset に設定する必要があります。ファイル名の形式例: |
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 セットの形式例: |
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 | この設定を有効にするには、 |
scan.only.deserialize.captured.tables.changelog.enabled | 増分フェーズで、指定されたテーブルの変更イベントのみを逆シリアル化するかどうかを指定します。 | いいえ | BOOLEAN |
| 有効な値:
|
scan.parallel-deserialize-changelog.enabled | 増分フェーズで、複数のスレッドを使用して変更イベントを解析するかどうかを指定します。 | いいえ | BOOLEAN | false | 有効な値:
説明 VVR 8.0.11 以降でのみサポートされています。 |
scan.parallel-deserialize-changelog.handler.size | 複数のスレッドを使用して変更イベントを解析する場合のイベントハンドラの数。 | いいえ | INTEGER | 2 | 説明 VVR 8.0.11 以降でのみサポートされています。 |
metadata-column.include-list | 下流に渡すメタデータ列。 | いいえ | STRING | なし | 利用可能なメタデータには、 説明 MySQL CDC YAML コネクタは、データベース名、テーブル名、および 重要
|
scan.newly-added-table.enabled | チェックポイントから再起動する際に、前回の起動時に一致しなかった新規追加テーブルを同期するか、または一致しなくなったテーブルを状態から削除するかを指定します。 | いいえ | BOOLEAN | false | これは、チェックポイントまたはセーブポイントから再起動するときに有効になります。 重要 完全データ読み取りフェーズ中に、セーブポイントを保存し、ソーステーブルに新しいテーブルを追加または削除してから、セーブポイントからジョブを再起動することはできません。これにより、ジョブがデータの読み取りに失敗します。 |
scan.binlog.newly-added-table.enabled | 増分フェーズで、一致する新規追加テーブルからのデータを送信するかどうかを指定します。 | いいえ | BOOLEAN | false |
|
scan.incremental.snapshot.chunk.key-column | スナップショットフェーズでのシャーディングの分割列として使用する特定のテーブルの列を指定します。 | いいえ | STRING | なし |
|
scan.parse.online.schema.changes.enabled | 増分フェーズで、RDS ロックレス変更 DDL イベントの解析を試みるかどうかを指定します。 | いいえ | BOOLEAN | false | 有効な値:
これは実験的な機能です。オンラインでのロックレス変更を実行する前に、回復のために Flink ジョブのスナップショットを取得してください。 説明 VVR 11.0 以降でのみサポートされています。 |
scan.incremental.snapshot.backfill.skip | スナップショット読み取りフェーズ中にバックフィルをスキップするかどうかを指定します。 | いいえ | BOOLEAN | 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 | 有効な値:
|
treat-timestamp-as-datetime-enabled | TIMESTAMP 型を DATETIME 型として扱うかどうかを指定します。 | いいえ | BOOLEAN | false | 有効な値:
MySQL の TIMESTAMP 型は UTC 時刻を格納し、タイムゾーンの影響を受けます。MySQL の DATETIME 型はリテラルな時刻を格納し、タイムゾーンの影響を受けません。 このパラメーターを有効にすると、MySQL の TIMESTAMP 型データが server-time-zone に基づいて DATETIME 型に変換されます。 |
include-comments.enabled | テーブルと列のコメントを同期するかどうかを指定します。 | いいえ | BOOELEAN | false | 有効な値:
このオプションを有効にすると、ジョブのメモリ使用量が増加します。 |
scan.incremental.snapshot.unbounded-chunk-first.enabled | スナップショット読み取りフェーズ中に、有界でないチャンクを最初にディスパッチするかどうかを指定します。 | いいえ | BOOELEAN | 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 | なし | これは、データ読み取りを制限する必要があるシナリオに適用されます。この制限は、完全フェーズと増分フェーズの両方で有効です。 ソースの 完全データ読み取りフェーズでは、通常、各バッチで読み取る行数を減らす必要があります。 説明 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)] |
|
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: 50debezium.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: trueMySQL CDC Enterprise Edition のバイナリログ消費能力は 85 MB/s で、これはオープンソースコミュニティ版の約 2 倍です。バイナリログファイルの生成速度が 85 MB/s (つまり、512 MB のファイルが 6 秒ごとに 1 つ) を超えると、Flink ジョブの遅延は増加し続けます。バイナリログファイルの生成速度が遅くなると、処理遅延は徐々に減少します。バイナリログファイルに大規模トランザクションが含まれている場合、処理遅延が一時的に増加することがあります。そのトランザクションのログが読み取られた後、遅延は減少します。
データ遅延の診断によるジョブスループットの最適化
増分フェーズでデータ遅延が発生した場合は、以下の手順で問題を分析してください:
概要ページの
currentFetchEventTimeLagおよびcurrentEmitEventTimeLagメトリックを確認します。currentFetchEventTimeLagメトリックは、バイナリログからのデータ読み取りの遅延を表します。currentEmitEventTimeLagメトリックは、ジョブに関連するテーブルのバイナリログからのデータ読み取りの遅延を表します。シナリオ
説明
currentFetchEventTimeLagは低いが、currentEmitEventTimeLagは高く、ほとんど更新されない。currentFetchEventTimeLagが低いことは、データベースからのバイナリログのプルが効率的であることを示します。しかし、バイナリログにはジョブが読み取る必要のあるテーブルのデータがほとんど含まれていません。そのため、currentEmitEventTimeLagはほとんど更新されません。これは期待される動作です。currentFetchEventTimeLagとcurrentEmitEventTimeLagの両方が高い。これは、ソーステーブルの読み取り性能が低いことを示します。このセクションの以降の手順に進んで最適化を行ってください。
バックプレッシャーは、ソースが下流オペレーターにデータを送信するレートを低下させる可能性があります。sourceIdleTime が定期的に増加し、currentFetchEventTimeLag と currentEmitEventTimeLag の両方が継続的に増加することが観察される場合があります。これを解決するには、バックプレッシャーが発生しているノードの並列度を上げてください。
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.binlog.newly-added-table.enabledを同時に有効にしないでください。両方を有効にすると、データが重複する可能性があります。