Flink CDC を使用して、Realtime Compute for Apache Flink の Iceberg REST を介してデータを Data Lake Formation (DLF) カタログに同期します。
前提条件
-
Realtime Compute for Apache Flink のフルマネージドワークスペースが必要です。作成していない場合は、「Realtime Compute for Apache Flink の有効化」をご参照ください。
-
Realtime Compute for Apache Flink ワークスペースと DLF が同じリージョンにあることを確認してください。また、ワークスペースの VPC を DLF ホワイトリストに追加する必要があります。詳細については、「VPC ホワイトリストの設定」をご参照ください。
制限事項
DLF への Iceberg REST 接続には、Realtime Compute for Apache Flink エンジンバージョン VVR 11.6.0 以降が必要です。
Flink での DLF カタログの登録
この操作により、DLF カタログへのマッピングが作成されます。Flink でカタログを作成または削除しても、DLF 内の実際のデータには影響しません。
Iceberg REST を介して DLF カタログに作成されたすべてのテーブルは、Iceberg テーブルです。
ワークスペースの [アクション] 列で、[コンソール] をクリックします。
左側のナビゲーションメニューで、 をクリックします。
新しいスクリプトを作成します。SQL エディターで、次の SQL 文をコピーして貼り付けます。右下隅で [環境] をクリックし、VVR 11.2.0 以降のセッションクラスターを選択して、SQL 文を実行し、Iceberg REST を介して DLF カタログを登録します。
CREATE CATALOG `catalog_name` WITH ( 'type' = 'iceberg', 'catalog-type' = 'rest', 'uri' = 'http://cn-hangzhou-vpc.dlf.aliyuncs.com/iceberg', 'warehouse' = 'iceberg_test', 'rest.signing-region' = 'cn-hangzhou', 'io-impl' = 'org.apache.iceberg.rest.DlfFileIO' );次の表にオプションを説明します。
オプション
説明
必須
例
type
タイプ。これを
icebergに設定します。はい
iceberg
catalog-type
カタログタイプ。これを
restに設定します。はい
rest
token.provider
トークンプロバイダー。これを dlf に設定します。
はい
dlf
uri
Iceberg REST を介して DLF カタログにアクセスするために使用される URI。詳細については、「Iceberg REST」をご参照ください。
はい
http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg
warehouse
DLF カタログの名前。
はい
iceberg_test
rest.signing-region
DLF のリージョン ID。詳細については、「エンドポイント」をご参照ください。
はい
ap-southeast-1
io-impl
これを
org.apache.iceberg.rest.DlfFileIOに設定します。はい
org.apache.iceberg.rest.DlfFileIO
カタログを使用するための Flink CDC の設定
データインジェストジョブを作成するには、「Flink CDC データインジェストジョブの開発」をご参照ください。
すでに Flink カタログマッピングを作成している場合は、「既存のカタログを再利用して接続情報を取得する」を参照し、データインジェストジョブに次のシンク設定を追加します。
sink:
type: iceberg
using.built-in-catalog: catalog_name
設定例
以下の例は、Flink CDC YAML ジョブを使用した一般的なデータ同期パターンを示しています。
MySQL データベース全体の DLF への同期
このジョブは、MySQL データベース全体を DLF に同期します。
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (オプション) 増分フェーズ中に作成されたテーブルを同期します。
scan.binlog.newly-added-table.enabled: true
# (オプション) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
# (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
MySQL ソースでは、これらのオプションパラメーターの使用を推奨します。詳細については、「MySQL」をご参照ください。
-
パラメーター:
scan.binlog.newly-added-table.enabled機能:増分フェーズ中に作成されたテーブルを同期します。
-
パラメーター:
include-comments.enabled機能:テーブルとフィールドのコメントを同期します。
-
パラメーター:
scan.incremental.snapshot.unbounded-chunk-first.enabled機能:TaskManager の OOM エラーを防ぎます。
-
パラメーター:
scan.only.deserialize.captured.tables.changelog.enabled機能:一致したテーブルのみを解析することで読み取りを高速化します。
パーティション分割された DLF テーブルへの書き込み
partition-keys パラメーターを使用してパーティションキーを指定します。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。例:
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (オプション) 増分フェーズ中に作成されたテーブルを同期します。
scan.binlog.newly-added-table.enabled: true
# (オプション) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
# (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (オプション) パーティションキーを設定します。
partition-keys: id,pt
- source-table: mysql_test.tbl2
partition-keys: id,pt
追記専用の DLF テーブルへの書き込み
データインジェストジョブは、ソースからのすべての変更イベントタイプをキャプチャします。DELETE 操作をダウンストリームの INSERT 操作に変換してソフトデリートを実装するには、以下のようにジョブを設定します。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (オプション) 増分フェーズ中に作成されたテーブルを同期します。
scan.binlog.newly-added-table.enabled: true
# (オプション) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
# (オプション) TaskManager の OOM エラーを防ぐために、非有界チャンクを優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (オプション) 一致したテーブルのみを解析することで読み取りを高速化します。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (オプション) パーティションキーを設定します。
partition-keys: id,pt
# (オプション) ソフトデリートを実装します。
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
- source-table: mysql_test.tbl2
# (オプション) パーティションキーを設定します。
partition-keys: id,pt
# (オプション) ソフトデリートを実装します。
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
-
プロジェクションに
__data_event_type__を追加すると、変更イベントタイプがダウンストリームテーブルに新しいフィールドとして書き込まれます。converter-after-transformをSOFT_DELETEに設定すると、DELETE 操作が INSERT に変換されます。これにより、ダウンストリームテーブルにはすべての変更イベントが記録されます。詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。
リアルタイム Kafka CDC データの DLF への同期
Kafka の inventory トピックに、customers と products の 2 つのテーブルの変更データが Debezium JSON 形式で保存されているとします。このジョブは、これらのテーブルのデータを対応する DLF ターゲットに同期します。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: debezium-json
debezium-json.distributed-tables: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Debezium JSON には主キー情報がないため、手動で追加します。
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
-
Kafka ソースは、
canal-json、debezium-json(デフォルト)、json形式をサポートしています。 -
debezium-jsonを使用する場合、Debezium JSON メッセージには主キー情報が含まれていないため、変換ルールを使用して手動で主キーを追加する必要があります。transform: - source-table: \.*.\.* projection: \* primary-keys: id -
単一テーブルのデータが複数のパーティションにまたがっている場合や、パーティションをまたいでテーブルをマージする必要がある場合は、
debezium-json.distributed-tablesまたはcanal-json.distributed-tablesをtrueに設定します。 -
Kafka ソースは、複数のスキーマ推論戦略をサポートしています。
schema.inference.strategyパラメーターを使用して、優先する戦略を設定します。詳細については、「Kafka」をご参照ください。
リアルタイム Kafka ログの DLF への同期
Kafka クラスターがカスタム JSON 形式でデータを保存している場合、Flink CDC YAML ジョブを設定して DLF に同期できます。システムは、データ型の推論、スキーマ推論、およびスキーマ進化を自動的に処理します。
inventory トピックに、単一のログテーブルのデータが JSON 形式で保存されているとします。このジョブは、そのデータを対応する DLF テーブルに同期します。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
json.infer-schema.flatten-nested-columns.enable: true
# (オプション) 最初の 100 件の解析エラーをスキップします。エラーが 100 件を超えるとジョブは失敗します。
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# テーブルに主キーを追加します。
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# inventory トピックのすべてのデータを test_database.inventory に書き込みます。
route:
- source-table: inventory
sink-table: test_database.inventory
pipeline:
# (オプション) 処理例外を引き起こしたダーティデータをログに記録します。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
inventory トピックに、複数のログテーブルのデータが JSON 形式で保存されており、データベース名とテーブル名が databaseName フィールドと tableName フィールドにあるとします。このジョブは、これらのテーブルのデータを対応する DLF ターゲットに同期します。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (オプション) JSON データ内のネストされた列を再帰的にフラット化します。
json.infer-schema.flatten-nested-columns.enable: true
# databaseName と tableName フィールドの値をデータベース名とテーブル名として使用します。
json.decode.parser-table-id.fields: databaseName,tableName
# (オプション) 最初の 100 件の解析エラーをスキップします。エラーが 100 件を超えるとジョブは失敗します。
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# テーブルに主キーを追加します。
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# ods.inventory、ods.customer、ods.user をそれぞれ test_database.inventory、test_database.customer、test_database.user に書き込みます。
route:
- source-table: ods.inventory
sink-table: test_database.inventory
- source-table: ods.customer
sink-table: test_database.customer
- source-table: ods.user
sink-table: test_database.user
pipeline:
# (オプション) 処理例外を引き起こしたダーティデータをログに記録します。
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
JSON 形式の Kafka ソーステーブルのスキーマ推論と進化戦略の詳細については、「スキーマ解析と変更同期戦略」をご参照ください。
ジョブ設定の詳細については、「Flink CDC データインジェストジョブ開発リファレンス」をご参照ください。