Realtime Compute for Apache Flink では、YAML ファイルを使用して Flink CDC データインジェストジョブを作成し、ソースからシンクへデータを同期できます。このトピックでは、Flink CDC データインジェストジョブを開発する手順について説明します。
背景情報
YAML 設定を使用すると、複雑な ETL パイプラインを簡単に定義でき、これらは自動的に Flink の実行ロジックに変換されます。Flink CDC データインジェストは、Flink CDC に基づく強力なデータ統合ソリューションを提供します。このソリューションは、データベース全体の同期、単一テーブルの同期、シャーディングされたデータベースとテーブルの同期、新しいテーブルの自動検出、スキーマ変更の処理、およびカスタム計算列を効率的にサポートします。また、ETL 処理、WHERE 句フィルタリング、およびカラムプルーニングもサポートしています。この宣言的なアプローチにより、データ統合プロセスが大幅に簡素化され、効率と信頼性が向上します。
Flink CDC の利点
Realtime Compute for Apache Flink では、Flink CDC データインジェストジョブ、SQL ジョブ、または DataStream ジョブを開発してデータ同期を実行できます。以下のセクションでは、他の 2 つの方法と比較した場合の Flink CDC データインジェストジョブの利点について説明します。
Flink CDC と Flink SQL の比較
Flink CDC データインジェストジョブと SQL ジョブは、データ転送に異なるデータ型を使用します。
-
SQL ジョブは、データを
RowDataとして送信します。各RowDataオブジェクトには独自の変更タイプがあり、主なタイプとして挿入 (+I)、更新前 (-U)、更新後 (+U)、削除 (-D) の 4 種類があります。 -
Flink CDC は
SchemaChangeEventを使用して、テーブルの作成、カラムの追加、テーブルの切り捨てなどのスキーマ変更情報を伝達します。また、DataChangeEventを使用して、挿入、更新、削除などのデータ変更を伝達します。更新メッセージには変更前と変更後の両方の内容が含まれているため、元の変更データをシンクに書き込むことができます。
次の表に、SQL ジョブに対する Flink CDC データインジェストジョブの利点を示します。
|
Flink CDC インジェスト |
Flink SQL |
|
スキーマの自動検出とデータベース全体の同期 |
手動の |
|
複数のスキーマ変更ポリシーをサポート |
スキーマ変更をサポートしない |
|
元のチェンジログを保持 |
元のチェンジログ構造が損なわれる |
|
複数テーブルの読み書きをサポート |
単一テーブルの読み書き |
CTAS/CDAS 文と比較して、Flink CDC ジョブは、以下のような、より強力な機能を提供します。
-
新しいデータの書き込みによる同期のトリガーを待つことなく、アップストリームのスキーマ変更を即座に同期。
-
元のチェンジログの保持、更新メッセージが分割されないことを保証。
-
TRUNCATE TABLEやDROP TABLEなど、より多くの種類のスキーマ変更の同期。 -
柔軟なテーブルマッピングとシンクテーブル名の定義。
-
柔軟で設定可能なスキーマ進化の動作。
-
WHERE句によるデータフィルタリング -
カラムプルーニングのサポート。
Flink CDC と Flink DataStream の比較
次の表に、DataStream ジョブに対する Flink CDC データインジェストジョブの利点を示します。
|
Flink CDC インジェスト |
Flink DataStream |
|
すべてのスキルレベルのユーザーが利用可能です。 |
Java と分散システムに関する専門知識が必要です。 |
|
基盤となる複雑さを隠蔽し、開発を簡素化します。 |
Flink フレームワークの知識が必要です。 |
|
習得しやすい YAML 形式。 |
依存関係管理のため、Maven のようなツールの知識が必要です。 |
|
既存ジョブの再利用性が高いです。 |
既存コードの再利用が困難です。 |
制限事項
-
Flink CDC データインジェストジョブを開発するには、Ververica Runtime (VVR) 11.1 以降を使用します。VVR 8.x を使用する必要がある場合は、VVR 8.0.11 を使用してください。
-
各ジョブは、1 つのソースと 1 つのシンクのみをサポートします。複数のソースから読み取ったり、複数のシンクに書き込んだりするには、複数の Flink CDC ジョブを作成する必要があります。
-
Flink CDC ジョブは、セッションクラスターにはデプロイできません。
-
Flink CDC データインジェストジョブは、自動チューニングをサポートしていません。
Flink CDC データインジェストコネクタ
Flink CDC データインジェストでサポートされているソースおよびシンクコネクタの詳細については、「サポートされているコネクタ」をご参照ください。
Flink CDC データインジェストジョブの作成
テンプレートから作成
-
対象のワークスペースの Actions 列で、 [Console] をクリックします。
-
左側メニューで、 を選択します。
-
をクリックし、 [New Draft with Template] をクリックします。 -
データ同期テンプレートを選択します。
現在、MySQL から StarRocks、MySQL から Paimon、および MySQL から Hologres へのテンプレートのみが利用可能です。
-
[name]、 [location]、 [engine version] などのジョブ情報を入力し、 [OK] をクリックします。
-
Flink CDC ジョブのソースとシンクの情報を設定します。
パラメータ設定の詳細は、関連するコネクタのドキュメントをご参照ください。
CTAS/CDAS ジョブから作成
-
ジョブに複数の CXAS 文が含まれている場合、Flink は最初の文のみを検出して変換します。
-
Flink SQL と Flink CDC では、組み込み関数のサポートに違いがあるため、生成された
transformルールはそのままでは機能しない場合があります。必要に応じて確認し、調整する必要があります。 -
ソースが MySQL で、元の CTAS/CDAS ジョブがまだ実行中の場合、競合を防ぐために、Flink CDC データ取り込みジョブでソースの
server-idを調整する必要があります。
-
対象のワークスペースの [Actions] 列で、 [Console] をクリックします。
-
左側メニューで、 を選択します。
-
をクリックし、 [New Draft from CTAS/CDAS] をクリックします。対象の CTAS または CDAS ジョブを選択し、 [OK] をクリックします。選択ページでは、有効な CTAS および CDAS ジョブのみが表示されます。通常の ETL ジョブや構文エラーのあるドラフトは表示されません。
-
[name]、 [location]、 [engine version] などのジョブ情報を入力し、 [OK] をクリックします。
オープンソース Flink CDC から作成
-
対象のワークスペースの [Actions] 列で、 [Console] をクリックします。
-
左側メニューで、 を選択します。
-
をクリックし、 [New Draft] を選択します。 [name] と [engine version] を入力し、 [Create] をクリックします。 -
オープンソースの Flink CDC ジョブの YAML コードをエディターに貼り付けます。
-
(任意) [Validate] をクリックします。
このオプションは、構文エラー、ネットワーク接続の問題、および権限の問題をチェックします。
ゼロから作成
-
対象のワークスペースの [Actions] 列で、 [Console] をクリックします。
-
左側メニューで、 を選択します。
-
をクリックし、 [New Draft] を選択します。 [name] と [engine version] を入力し、 [Create] をクリックします。 -
YAML を使用して Flink CDC ジョブを設定します。例:
# 必須 source: # データソースのタイプ。 type: <ソースコネクタのタイプに置き換えてください> # データソースの設定。設定項目の詳細は、対応するコネクタのドキュメントをご参照ください。 ... # 必須 sink: # シンクのタイプ。 type: <シンクコネクタのタイプに置き換えてください> # シンクの設定。設定項目の詳細は、対応するコネクタのドキュメントをご参照ください。 ... # 任意 transform: # flink_test.customers テーブルの transform ルール。 - source-table: flink_test.customers # プロジェクション設定。同期するカラムを指定して、データ変換を実行します。 projection: id, username, UPPER(username) as username1, age, (age + 1) as age1, test_col1, __schema_name__ || '.' || __table_name__ identifier_name # フィルター条件。id が 10 より大きいデータのみを同期します。 filter: id > 10 # transform ルールの説明。 description: append calculated columns based on source table # 任意 route: # ソーステーブルとシンクテーブルのマッピングを指定する route ルール。 - source-table: flink_test.customers sink-table: db.customers_o # route ルールの説明。 description: sync customers table - source-table: flink_test.customers_suffix sink-table: db.customers_s # route ルールの説明。 description: sync customers_suffix table # 任意 pipeline: # ジョブの名前。 name: MySQL to Hologres Pipeline説明Flink CDC ジョブでは、キーと値をコロンとスペースで区切ります。形式は
Key: Valueです。次の表に、コードブロックについて説明します。
必須
モジュール
説明
はい
source
データパイプラインの開始点。Flink CDC は、ソースから変更データをキャプチャします。
説明-
現在、サポートされているソースは MySQL のみです。設定の詳細については、「MySQL コネクタ」をご参照ください。
-
変数を使用して機密情報を管理できます。詳細については、「変数管理」をご参照ください。
sink
データパイプラインのエンドポイント。Flink CDC は、キャプチャしたデータ変更をシンクシステムに転送します。
説明-
サポートされているシンクについては、「Flink CDC データインジェストコネクタ」をご参照ください。設定の詳細については、特定のコネクタのドキュメントをご参照ください。
-
変数を使用して機密情報を管理できます。詳細については、「変数管理」をご参照ください。
いいえ
pipeline
(データパイプライン)
パイプライン名など、データパイプラインジョブ全体の基本設定を定義します。
transform
Flink パイプラインを流れるデータを操作するためのデータ変換ルールを指定します。ETL 処理、
WHERE句フィルタリング、列プルーニング、および計算列をサポートします。Flink CDC によってキャプチャされた生の変更データを、特定のダウンストリームシステムに合わせて変換するには、
transformブロックを使用します。route
このモジュールが設定されていない場合、ジョブはデフォルトでデータベース全体または対象テーブルの同期を実行します。
場合によっては、特定のルールに基づいて、キャプチャされた変更データを異なる宛先に送信する必要があります。
routeモジュールでは、ソースとシンクのマッピング関係を柔軟に指定して、データを異なるターゲットに送信できます。各モジュールの構文と設定の詳細は、「Flink CDC データインジェストジョブのリファレンス」をご参照ください。
次のコードは、MySQL の
app_dbデータベースのすべてのテーブルを Hologres のデータベースに同期する例を示しています。source: type: mysql hostname: <ホスト名> port: 3306 username: ${secret_values.mysqlusername} password: ${secret_values.mysqlpassword} tables: app_db.\.* server-id: 5400-5404 # (任意) 増分フェーズで新しく作成されたテーブルからデータを同期します。 scan.binlog.newly-added-table.enabled: true # (任意) テーブルとフィールドのコメントを同期します。 include-comments.enabled: true # (任意) 無制限のチャンクを優先的にディスパッチして、潜在的な TaskManager の OutOfMemory 問題を防ぎます。 scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (任意) 読み取りを高速化するため、解析フィルターを有効にします。 scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: hologres name: Hologres Sink endpoint: <エンドポイント> dbname: <データベース名> username: ${secret_values.holousername} password: ${secret_values.holopassword} pipeline: name: Sync MySQL Database to Hologres -
-
(任意) [Validate] をクリックします。
このオプションは、構文エラー、ネットワーク接続の問題、および権限の問題をチェックします。
関連トピック
-
Flink CDC ジョブの開発が完了したら、デプロイする必要があります。詳細については、「ジョブのデプロイ」をご参照ください。
-
MySQL データベースから StarRocks にデータを同期する Flink CDC ジョブを迅速に構築する方法については、「チュートリアル:Flink CDC データインジェストジョブの作成」をご参照ください。