このトピックでは、ストリーミングデータレイクハウス用の Paimon コネクタの使用方法について説明します。最良の結果を得るには、コネクタを Paimon カタログと併用することを推奨します。
背景情報
Apache Paimon は、高スループットの書き込みと低レイテンシーのクエリをサポートする、ストリーミングとバッチを統合したデータレイクストレージフォーマットです。Paimon は、Alibaba Cloud E-MapReduce で利用可能な Flink、Spark、Hive、Trino などの一般的なコンピュートエンジンと密接に統合されています。Apache Paimon を使用すると、HDFS または OSS 上にデータレイクを迅速に構築し、これらのコンピュートエンジンに接続してデータレイク分析を行うことができます。詳細については、Apache Paimon をご参照ください。
|
カテゴリ |
説明 |
|
サポートされるタイプ |
ソーステーブル、ディメンションテーブル、結果テーブル、およびデータインジェストのターゲット |
|
実行モード |
ストリーミングモードとバッチモード |
|
データフォーマット |
サポートされていません |
|
監視メトリック |
なし |
|
API タイプ |
データインジェスト用の SQL と YAML |
|
結果テーブルの更新と削除 |
はい |
主な特徴
Apache Paimon は、以下の主要な機能を提供します:
-
HDFS またはオブジェクトストレージ上に、軽量で低コストのデータレイクを構築します。
-
ストリーミングモードとバッチモードの両方で、大規模なデータセットの読み取りと書き込みを行います。
-
数分から数秒のデータ鮮度で、バッチおよび OLAP クエリを実行します。
-
増分データをインジェストおよび生成し、従来のオフラインおよび最新のストリーミングデータウェアハウスの両方のストレージレイヤーとして機能します。
-
データを事前集計して、ストレージコストと下流の計算負荷を削減します。
-
データの履歴バージョンにアクセスします。
-
データを効率的にフィルターします。
-
スキーマ進化を有効にします。
制限事項と推奨事項
-
Paimon コネクタには、Flink コンピュートエンジン VVR 6.0.6 以降が必要です。
-
次の表に、Paimon と VVR のバージョン互換性を示します。
Apache Paimon バージョン
VVR
1.3.1
11.5、11.6、11.7、11.8
1.3
11.4
1.2
11.2、11.3
1.1
11.1
1.0
8.0.11
-
同時書き込みに関するストレージの推奨事項
複数のジョブが同じ Paimon テーブルに同時に書き込みを行う場合、標準の OSS ストレージ (oss://) を使用すると、アトミックなファイル操作の制限により、コミットの競合やジョブの失敗が時折発生する可能性があります。
安定した一貫性のある書き込みを行うには、強力なアトミック保証を提供するメタデータまたはストレージサービスを使用してください。推奨されるオプションは、Paimon のメタデータとストレージの統合管理を提供する Data Lake Formation (DLF) です。あるいは、OSS-HDFS または HDFS を使用することもできます。
-
設定変更の反映方法
Paimon テーブルの設定パラメーターへの変更は、関連するジョブを再起動した後にのみ有効になります。実行中のジョブは、これらの変更を動的に読み込みません。
-
削除されたパーティションの物理的な回収の遅延
DROP PARTITION 操作を実行しても、システムは基になる物理データファイルを即座に削除しません。
この操作は論理削除を実行します。Paimon は、最新のスナップショットからのみターゲットパーティションのメタデータを削除します。Paimon はタイムトラベル機能をサポートしているため、履歴スナップショットは依然としてそのパーティションのデータファイルを参照しています。物理データファイルは、パーティションを参照するすべての履歴スナップショットが保持期間の上限に達し、スナップショットの有効期限切れメカニズムによってクリーンアップされた後にのみ、完全に削除されます。
SQL
SQL ジョブで Paimon コネクタをソーステーブルまたは結果テーブルとして使用します。
構文
-
Paimon カタログで Paimon テーブルを作成する場合、
connectorパラメーターを指定する必要はありません。構文は次のとおりです:CREATE TABLE `<YOUR-PAIMON-CATALOG>`.`<YOUR-DB>`.paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( ... );説明Paimon カタログに Paimon テーブルを既に作成している場合は、直接使用できます。
-
別のカタログに Paimon 一時テーブルを作成する場合は、'connector' と 'path' パラメーターを指定する必要があります。構文は次のとおりです:
CREATE TEMPORARY TABLE paimon_table ( id BIGINT, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'paimon', 'path' = '<path-to-paimon-table-files>', 'auto-create' = 'true', -- 指定されたパスに Paimon テーブルのデータファイルが存在しない場合、自動的に作成されます。 ... );説明-
パスの例:
'path' = 'oss://<bucket>/test/order.db/orders'。.dbサフィックスを省略しないでください。Paimon はこのサフィックスに依存してデータベースを識別します。 -
同じテーブルに書き込む複数のジョブは、同じパス設定を共有する必要があります。
-
2つのパス設定が異なる場合、Paimon はそれらを同じテーブルとして認識しません。物理パスが同じであっても、一貫性のないカタログ設定は、同時書き込みの競合、コンパクション操作の失敗、およびデータ損失につながる可能性があります。たとえば、Paimon は
oss://b/testとoss://b/test/を、末尾のスラッシュのために異なるテーブルと見なします。これらが同じ物理的な場所を指している場合でも同様です。
-
WITH パラメーター
|
パラメーター |
説明 |
タイプ |
必須 |
デフォルト |
備考 |
|
connector |
テーブルのコネクタを指定します。 |
文字列 |
いいえ |
なし |
|
|
path |
テーブルのストレージパス。 |
文字列 |
いいえ |
なし |
|
|
auto-create |
指定されたパスにテーブルファイルが存在しない場合に、自動的に作成するかどうかを指定します。 |
ブール値 |
いいえ |
false |
有効な値:
|
|
file.format |
データファイルのフォーマット。 |
文字列 |
いいえ |
parquet |
有効な値:
|
|
bucket |
パーティションごとのバケット数。 |
整数 |
いいえ |
1 |
Paimon は 説明
各バケットには 5 GB 未満のデータを含めることを推奨します。 |
|
bucket-key |
バケットキーとして使用される列。 |
文字列 |
いいえ |
なし |
データをバケットに分散するために使用される列を指定します。 複数の列名はカンマ (,) で区切ります。例: 説明
|
|
changelog-producer |
Changelog の生成メカニズム。 |
文字列 |
いいえ |
none |
Paimon は任意の入力ストリームに対して完全な Changelog を生成できます。これは、すべての
Changelog プロデューサーの選択方法の詳細については、Changelog の生成をご参照ください。 |
|
full-compaction.delta-commits |
2つの連続するフルコンパクション間のコミットの最大数。 |
整数 |
いいえ |
なし |
フルコンパクションをトリガーする前に許可されるスナップショットコミットの最大数を指定します。 |
|
lookup.cache-max-memory-size |
Paimon ディメンションテーブルのメモリキャッシュサイズ。 |
文字列 |
いいえ |
256 MB |
このパラメーターは、ディメンションテーブルのルックアップと |
|
merge-engine |
同じプライマリキーを持つレコードをマージするメカニズム。 |
文字列 |
いいえ |
deduplicate |
有効な値:
マージエンジンの詳細な分析については、マージエンジンをご参照ください。 |
|
partial-update.ignore-delete |
削除 (-D) メッセージを無視するかどうかを指定します。 |
ブール値 |
いいえ |
false |
有効な値:
説明
|
|
ignore-delete |
削除 (-D) メッセージを無視するかどうかを指定します。 |
ブール値 |
いいえ |
false |
有効な値は partial-update.ignore-delete のものと同じです。 説明
|
|
partition.default-name |
デフォルトのパーティション名。 |
文字列 |
いいえ |
__DEFAULT_PARTITION__ |
パーティション列の値が null または空の文字列の場合に使用するパーティション名。 |
|
partition.expiration-check-interval |
システムが期限切れのパーティションをチェックする頻度。 |
文字列 |
いいえ |
1h |
詳細については、自動パーティション有効期限の設定方法をご参照ください。 |
|
partition.expiration-time |
パーティションが期限切れになるまでの期間。 |
文字列 |
いいえ |
なし |
パーティションの経過時間がこの値を超えると、パーティションは期限切れになります。デフォルトでは、パーティションは期限切れになりません。 システムはパーティションの値からパーティションの経過時間を計算します。詳細については、自動パーティション有効期限の設定方法をご参照ください。 |
|
partition.timestamp-formatter |
時間文字列をタイムスタンプに変換するためのフォーマット文字列。 |
文字列 |
いいえ |
なし |
パーティション値からパーティションの経過時間を抽出するためのフォーマットを指定します。詳細については、自動パーティション有効期限の設定方法をご参照ください。 |
|
partition.timestamp-pattern |
パーティション値を時間文字列に変換するためのフォーマット文字列。 |
文字列 |
いいえ |
なし |
パーティション値から時間文字列を抽出するためのパターンを指定します。詳細については、自動パーティション有効期限の設定方法をご参照ください。 |
|
scan.bounded.watermark |
スキャンの終了を示す Watermark 値。ソーステーブルは、その Watermark がこの値を超えるとデータの生成を停止します。 |
Long |
いいえ |
なし |
N/A |
|
scan.mode |
Paimon ソーステーブルの消費位置を指定します。 |
文字列 |
いいえ |
default |
詳細については、Paimon ソーステーブルの消費位置の設定方法をご参照ください。 |
|
scan.snapshot-id |
Paimon ソーステーブルが消費を開始するスナップショットを指定します。 |
整数 |
いいえ |
なし |
詳細については、Paimon ソーステーブルの消費位置の設定方法をご参照ください。 |
|
scan.timestamp-millis |
Paimon ソーステーブルが消費を開始する時点を指定します。 |
整数 |
いいえ |
なし |
詳細については、Paimon ソーステーブルの消費位置の設定方法をご参照ください。 |
|
snapshot.num-retained.max |
保持する最近のスナップショットの最大数。 |
整数 |
いいえ |
2147483647 |
この条件または |
|
snapshot.num-retained.min |
保持する最近のスナップショットの最小数。 |
整数 |
いいえ |
10 |
N/A |
|
snapshot.time-retained |
スナップショットの保持期間。 |
文字列 |
いいえ |
1h |
この条件または |
|
write-mode |
Paimon テーブルの書き込みモード。 |
文字列 |
いいえ |
change-log |
有効な値:
書き込みモードの詳細については、書き込みモードをご参照ください。 |
|
scan.infer-parallelism |
Paimon ソーステーブルの並列度を自動的に推論するかどうかを指定します。 |
ブール値 |
いいえ |
true |
有効な値:
|
|
scan.parallelism |
Paimon ソーステーブルの並列度。 |
整数 |
いいえ |
なし |
説明
ジョブの タブで [リソースモード] が [モード] に設定されている場合、このパラメーターは無視されます。 |
|
sink.parallelism |
Paimon sink テーブルの並列度。 |
整数 |
いいえ |
なし |
説明
ジョブの タブで [リソースモード] が [モード] に設定されている場合、このパラメーターは無視されます。 |
|
sink.clustering.by-columns |
Paimon sink テーブルへの書き込みのためのクラスタリング列を指定します。 |
文字列 |
いいえ |
なし |
Paimon の append-only テーブル (プライマリキーのないテーブル) の場合、このパラメーターはバッチジョブでのクラスタ化された書き込みを有効にします。このプロセスは、指定された列でデータをグループ化することにより、クエリのパフォーマンスを向上させます。 複数の列名はカンマ (,) で区切ります。例: クラスタリングの詳細については、Apache Paimon 公式ドキュメントをご参照ください。 |
|
sink.delete-strategy |
システムがリトラクションメッセージ (-D/-U) を正しく処理することを保証するための検証戦略を指定します。 |
Enum |
いいえ |
NONE |
有効な値と、リトラクションメッセージを処理する際の sink 演算子の期待される動作:
説明
|
|
blob-as-descriptor |
Blob 列が読み取られたときに Blob 記述子バイトを出力するかどうかを指定します。 |
ブール値 |
いいえ |
false |
このパラメーターを true に設定すると、クエリは実際の Blob コンテンツの代わりに BlobDescriptor のシリアル化されたバイトを返します。このパラメーターは Blob 署名付き URL 機能と一緒に使用します。テーブル作成時にこのパラメーターを設定する必要はありません。SQL ヒントを使用して読み取り時に動的に設定できます。このパラメーターは VVR 11.9-preview1 以降でサポートされています。 説明
VVR 11.9-preview1 以降でのみサポートされています。 |
設定オプションの詳細については、Apache Paimon 公式ドキュメントをご参照ください。
ベクターテーブルのパラメーター
以下のパラメーターは VVR 11.8 以降でのみサポートされています。
|
パラメーター |
説明 |
データの型 |
必須 |
デフォルト値 |
備考 |
|
|
検索中にプローブする IVF クラスターの数。 |
整数 |
いいえ |
16 |
値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。 |
|
|
top_k × refine_factor の IVF 候補を取得し、Paimon テーブルに保存されている元のベクターを使用して再ランキングします。 |
Float |
いいえ |
無効 |
すべての IVF バリアントでデフォルトで無効になっています。レイテンシーよりも再現率が優先されるシナリオで、圧縮インデックス (ivf-pq や ivf-hnsw-sq など) に最適です。 |
|
|
検索中の HNSW 検索幅。 |
整数 |
いいえ |
0 |
値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。0 はネイティブライブラリのデフォルトが使用されることを示します。 |
|
|
Lumina DiskANN の検索リストサイズ。 |
整数 |
いいえ |
max(1.5 × top_k, 16) |
値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。 |
|
|
Lumina DiskANN の検索ビーム幅。 |
整数 |
いいえ |
4 |
— |
|
|
並列 Lumina 検索の数。 |
整数 |
いいえ |
5 |
— |
特徴の詳細
データの鮮度と一貫性
Paimon sink テーブルは、各 Flink ジョブのチェックポイント中にデータをコミットするために2フェーズコミットプロトコルを使用します。したがって、データの鮮度は Flink ジョブのチェックポイント間隔によって決まります。各コミットは最大2つのスナップショットを生成します。
2つの Flink ジョブが同じ Paimon テーブルに同時に書き込む場合、ジョブが異なるバケットに書き込むと、直列化可能な一貫性を達成します。ジョブが同じバケットに書き込む場合、スナップショット分離のみを達成します。これは、テーブルのデータが両方のジョブの結果の混合になる可能性があることを意味しますが、データ損失は発生しません。
マージエンジン
Paimon sink テーブルが同じプライマリキーを持つ複数のレコードを受け取ると、一意性を維持するためにそれらを1つのレコードにマージします。この動作は merge-engine パラメーターを設定することで制御できます。次の表に、利用可能なマージエンジンを説明します。
|
マージエンジン |
説明 |
|
Deduplicate |
重複排除エンジンがデフォルトです。同じプライマリキーを持つ複数のレコードについて、Paimon sink テーブルは最新のレコードのみを保持し、他は破棄します。 説明
最新のレコードが削除メッセージの場合、そのプライマリキーを持つすべてのレコードが破棄されます。 |
|
Partial Update |
部分更新エンジンを使用すると、複数のメッセージで増分的に更新することで完全なレコードを構築できます。同じプライマリキーを持つ新しいレコードが到着すると、その null でない値が既存のレコードの対応するフィールドを上書きします。エンジンは新しいレコードで null のフィールドを無視し、既存の値を保持します。 たとえば、Paimon sink テーブルが次の3つのレコードを順に受け取るとします:
最初の列がプライマリキーの場合、最終的にマージされたレコードは <1, 25.2, 10, 'This is a book'> になります。 説明
|
|
Aggregation |
一部のユースケースでは、レコードの集計値のみが必要な場合があります。集約エンジンは、指定した集計関数を使用して同じプライマリキーを持つレコードを結合します。各非プライマリキー列について、
説明
|
Changelog プロデューサー
changelog-producer パラメーターを設定して、Paimon が任意の入力ストリームに対して完全な Changelog (すべての update_after レコードに対応する update_before レコードがある) を生成するように設定します。次の表に、利用可能な Changelog プロデューサーを説明します。詳細については、Apache Paimon 公式ドキュメントをご参照ください。
|
プロデューサー |
説明 |
|
None |
たとえば、下流のコンシューマーが列の合計を計算する必要があり、最新の値が 5 であることしかわからない場合、合計をどのように更新するかを判断できません。以前の値が 4 であれば、合計は 1 増加し、以前の値が 6 であれば、合計は 1 減少します。 説明
データベースなどの下流のコンシューマーが |
|
Input |
このプロデューサーは、入力ストリーム自体が既に完全な Changelog である場合 (Change Data Capture (CDC) からのデータなど) にのみ使用してください。 |
|
Lookup |
高いデータ鮮度 (例:分単位) を必要とするユースケースには、このオプションを使用してください。 |
|
Full Compaction |
データ鮮度の要件が低いユースケース (例:時間単位) には、このオプションを使用してください。 |
書き込みモード
Paimon テーブルは、以下の書き込みモードをサポートしています。
|
モード |
説明 |
|
Change-log |
|
|
Append-only |
|
Flink CDC データインジェストのターゲット
Paimon テーブルは、単一テーブルまたはデータベース全体からのデータのリアルタイム同期をサポートしています。上流のスキーマ変更は、リアルタイムで Paimon テーブルに同期されます。詳細については、データインジェストセクションをご参照ください。
Variant 読み取りプルーニング
クエリによって参照される Variant フィールドのみを読み取り、I/O とメモリのオーバーヘッドを削減します。
有効化の方法
|
パラメーター |
デフォルト値 |
説明 |
|
|
false |
|
前提条件
-
Variant データは VVR 11.6 以降でのみ書き込み可能です。以前のバージョンや他のエンジンで書き込まれた Variant データは、読み取りプルーニングをサポートしません。
-
上流のデータが有効になるには、
variant.inferShreddingSchemaオプションを有効にして自動 Variant シュレッディング推論を有効にするか、variant.shreddingSchemaおよび関連オプションを明示的に設定して、どの Variant メンバーがシュレッディングに適しているかを指定する必要があります。詳細については、Coreoptions をご参照ください。
サポートされるシナリオ
読み取りプルーニングは、SQL が文字列キーで Variant フィールドにアクセスする場合に有効になります。例:
SELECT v['a'] FROM t; -- 単一レベル
SELECT v['a']['b'] FROM t; -- ネスト
SELECT v['a'], v['b'] FROM t; -- 複数フィールド
SELECT id, v['a'] FROM t; -- 通常の列との混合
SELECT v['a'] + 1 FROM t; -- 式で参照
サポートされないシナリオ
-
SELECT v FROM tのように、Variant 全体を直接参照する場合。 -
SELECT v, v['a'] FROM tのように、直接参照とフィールドアクセスを組み合わせる場合。 -
v[0]のような配列インデックスアクセス。 -
テーブルスキーマにネストされた Row フィールドが含まれている場合。この場合、読み取りプルーニングはテーブル内のどの Variant フィールドにも適用されません。
Blob 署名付き URL
Paimon は、OSS に保存されている Blob 列データに対して、短期間有効なパブリック HTTPS GET 署名付き URL を生成できます。マルチモーダルモデルサービスなどの外部サービスは、Flink ジョブ内でファイルバイトを転送することなく、この URL を使用してファイルコンテンツを直接取得できます。
使用制限
-
この機能は Flink コンピュートエンジン VVR 11.9-preview1 以降でのみサポートされています。
-
OSS に保存されている Paimon テーブルのみがサポートされます。メタストアのタイプは問いません。Filesystem、DLF、その他のメタストアタイプはすべてサポートされます。
-
生成される URL はパブリック HTTPS URL であり、パブリックネットワークアクセスを持つ外部サービスによって消費される必要があります。カタログが内部エンドポイントを使用している場合、外部サービスが生成された URL にアクセスできるように、パブリック HTTPS エンドポイントに変更することを推奨します。
-
通常の Blob 列を読み取るには、まず
blob-as-descriptorオプションを使用して BlobDescriptor バイトを取得し、そのバイトを関数に渡す必要があります。
構文
descriptor_to_presigned_url(source_table, descriptor, validity)
try_descriptor_to_presigned_url(source_table, descriptor, validity)
パラメーター
|
パラメーター |
タイプ |
説明 |
|
|
STRING リテラル |
Blob 列を含む Paimon テーブル。フォーマットは |
|
|
BYTES |
BlobDescriptor バイト。 |
|
|
INTERVAL |
URL の有効期間。値は正の秒数でなければなりません。例: |
戻り値
|
タイプ |
説明 |
|
STRING |
短期間有効なパブリック HTTPS GET 署名付き URL。 |
例
-- SQL ヒントを使用して Blob 列を記述子として読み取り、署名付き URL を生成する
SELECT sys.descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-- フォールトトレラントモード:行レベルのエラーで NULL を返し、他の動作は変更なし
SELECT sys.try_descriptor_to_presigned_url(
'default.image_table',
image,
INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
-
descriptor_to_presigned_urlは行レベルのエラーで例外をスローします。try_descriptor_to_presigned_urlは行レベルのエラーで NULL を返します。 -
生成された URL にはファイル拡張子が含まれず、URL が指すバイトは Blob データと同一です。下流のサービスは、URL のサフィックスではなく、返されたバイトに基づいてファイル形式を識別する必要があります。
-
署名付き URL は一時的なアクセス認証情報に相当します。署名付き URL が生成されたら、すぐに下流のサービスに送信してください。ログや結果テーブルに書き込んだり、永続化したりしないでください。
Paimon をディメンションテーブルとして使用する
Paimon テーブルはディメンションテーブルとして使用できます。JOIN 構文については、ディメンションテーブルの JOIN 文をご参照ください。
デフォルトでは、ルックアップは各並列インスタンスですべてのデータをロードします。このアプローチは、小さなディメンションテーブルにのみ適しています。大きなディメンションテーブルの場合は、以下で説明するシャッフルルックアップソリューションを使用してください。
パーティション化されたディメンションテーブル
ディメンションテーブルがパーティション化されており、最新の1つまたは2つのパーティションのデータのみが必要な場合は、動的パーティション読み込み機能を使用できます:
SELECT * FROM T
JOIN DIM /*+ OPTIONS('lookup.dynamic-partition'='max_pt()', 'lookup.dynamic-partition.refresh-interval'='1 h') */
FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col = D.col;
|
パラメーター |
データの型 |
デフォルト値 |
説明 |
|
lookup.dynamic-partition |
文字列 |
N/A |
|
|
lookup.dynamic-partition.refresh-interval |
Duration |
1 h |
システムがディメンションテーブルのパーティション更新をチェックする間隔。 |
大きなディメンションテーブル:固定バケットテーブル
VVR 8.0.8 以降でのみサポートされています。固定バケットテーブル (bucket > 0) の場合、シャッフルルックアップを使用して、バケットキーによってデータを並列インスタンスに分散できます。これにより、各インスタンスは割り当てられたバケットのデータのみをロードします:
SELECT /*+ LOOKUP('table'='D', 'shuffle'='true') */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-
結合キーはバケットキーでなければなりません。バケットキーはデフォルトでプライマリキーになります。
-
この機能は固定バケットテーブル (bucket > 0) のみでサポートされます。
大きなディメンションテーブル:非固定バケットテーブル
VVR 8.0.10 以降でのみサポートされています。動的バケットテーブルまたは追加テーブルの場合、SHUFFLE_HASH または REPLICATED_SHUFFLE_HASH を使用して、各並列インスタンスがすべてのデータを読み取り、必要な部分のみを保持するようにできます:
-- Shuffle Hash
SELECT /*+ SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
-- Replicated Shuffle Hash
SELECT /*+ REPLICATED_SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
SHUFFLE_HASH と REPLICATED_SHUFFLE_HASH の詳細については、ディメンションテーブルの JOIN 文をご参照ください。
データインジェスト
データインジェスト YAML ジョブで Paimon コネクタを sink として使用できます。
構文
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: filesystem
catalog.properties.warehouse: /path/warehouse
パラメーター
|
パラメーター |
説明 |
必須 |
タイプ |
デフォルト |
注 |
|
type |
コネクタのタイプ。 |
はい |
STRING |
なし |
値は |
|
name |
sink の名前。 |
いいえ |
STRING |
なし |
|
|
catalog.properties.metastore |
Paimon カタログのタイプ。 |
いいえ |
STRING |
filesystem |
有効な値:
|
|
catalog.properties.* |
Paimon カタログを作成するためのパラメーター。 |
いいえ |
STRING |
なし |
詳細については、Paimon カタログの管理をご参照ください。 |
|
table.properties.* |
Paimon テーブルを作成するためのパラメーター。 |
いいえ |
STRING |
なし |
詳細については、Paimon テーブルオプションをご参照ください。 |
|
table.properties.file.format |
テーブル内のデータファイルのストレージフォーマット。 |
String |
いいえ |
parquet |
有効な値:
|
|
catalog.properties.warehouse |
ファイルストレージのルートディレクトリ。 |
いいえ |
STRING |
なし |
このパラメーターは、 |
|
commit.user-prefix |
データファイルをコミットするためのユーザー名プレフィックス。 |
いいえ |
STRING |
なし |
説明
異なるジョブには異なるユーザー名を設定することを推奨します。これにより、コミットの競合を引き起こしたジョブを特定しやすくなります。 |
|
partition.key |
パーティションテーブルのパーティションキー。 |
いいえ |
STRING |
なし |
異なるテーブルは |
|
sink.cross-partition-upsert.tables |
プライマリキーにすべてのパーティションキーが含まれていない、クロスパーティションアップサートが必要なテーブルをリストします。 |
いいえ |
STRING |
なし |
クロスパーティション更新があるテーブルに適用されます。
重要
|
|
sink.commit.parallelism |
Commit 演算子の並列度を指定します。 |
いいえ |
INTEGER |
なし |
Commit 演算子がボトルネックになっている場合、このパラメーターを使用して並列度を上げ、パフォーマンスを向上させます。 このパラメーターは、Realtime Compute for Apache Flink 11.6 以降でのみサポートされています。 説明
このパラメーターを設定すると、演算子の並列度が変更されます。ステートフルなジョブを再起動する際には、ジョブが部分的な演算子状態を無視できるように |
|
sink.existing-table.schema-check.enabled |
既存のテーブルスキーマの互換性をチェックし、事前に作成されたテーブルにデータが書き込まれるときにスキーマを自動的に適応させるかどうかを指定します。 |
ブール値 |
いいえ |
true |
次のリストは、有効な値と、事前に作成されたテーブルにデータが書き込まれるときの動作を説明しています:
説明
このパラメーターは、Realtime Compute for Apache Flink VVR 11.8 以降でのみサポートされています。 |
既存のカタログの再利用
Realtime Compute for Apache Flink 11.5 以降では、Flink CDC データインジェストジョブでデータ管理ページから組み込みの Paimon カタログを直接参照できます。これにより、手動での設定が削減されます。
sink:
type: paimon
using.built-in-catalog: paimon_dlf_catalog
catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
データインジェストジョブは、Paimon カタログ内のすべてのパラメーターを自動的に再利用できます。これは、YAML ジョブで catalog.properties. プレフィックスを持つパラメーターを手動で設定するのと同等です。
自動的に再利用されるパラメーターをオーバーライドするには、YAML ジョブで明示的に設定します。明示的な YAML 設定の方が優先度が高くなります。たとえば、前述のサンプルでは、fs.oss.endpoint パラメーターは YAML ジョブの値を使用し、paimon_dlf_catalog の値をオーバーライドします。
スキーマ変更
CDC YAML パイプラインジョブは、異なる戦略を使用してスキーマ変更を処理できます。戦略を設定するには、スキーマ進化の設定で説明されているように、パイプラインレベルのパラメーター schema.change.behavior を設定します。schema.change.behavior は IGNORE、LENIENT、TRY_EVOLVE、EVOLVE、または EXCEPTION に設定できます。次のセクションでは、スキーマ変更を伴う LENIENT および EVOLVE モードが、さまざまなスキーマ変更イベントをどのように処理するかを説明します。
LENIENT (デフォルト)
LENIENT モードでは、以下のスキーマ変更がサポートされています:
-
nullable 列の追加:列は結果テーブルのスキーマに自動的に追加され、新しい列のデータは自動的に同期されます。
-
nullable 列の削除:列は結果テーブルから削除されません。代わりに、列のデータは自動的に NULL 値に置き換えられます。
-
NOT NULL 列の追加:列は結果テーブルのスキーマに自動的に追加され、新しい列のデータは自動的に同期されます。新しい列はデフォルトで nullable であり、列が追加される前に生成されたデータは NULL に設定されます。
-
列名の変更:操作は列の追加と列の削除として扱われます。名前が変更された列は結果テーブルの末尾に追加され、元の列のデータは自動的に NULL 値に置き換えられます。たとえば、col_a が col_b に名前変更された場合、col_b は結果テーブルの末尾に追加され、col_a のデータは NULL 値に置き換えられます。
-
列タイプの変更:結果テーブルの対応する列のデータ型を変更します。
-
以下のスキーマ変更はサポートされていません:
-
プライマリキータイプの変更。
-
NULLABLE から NOT NULL への変更。
-
EVOLVE
EVOLVE モードでは、以下のスキーマ変更がサポートされています:
-
nullable 列の追加:サポートされています。
-
nullable 列の削除:サポートされていません。
-
NOT NULL 列の追加:結果テーブルに NULL 可能列が追加されます。
-
列名の変更:サポートされています。元の列は結果テーブルで名前が変更されます。
-
列タイプの変更:結果テーブルの対応する列のデータ型を変更します。
-
以下のスキーマ変更はサポートされていません:
-
プライマリキータイプの変更。
-
NULLABLE から NOT NULL への変更。
-
EVOLVE モードを有効にする例については、EVOLVE モードの有効化をご参照ください。
下流の Paimon テーブルが既に存在する場合、ジョブは既存のテーブルスキーマにデータを書き込み、テーブルを再度作成しようとはしません。既存のテーブルスキーマが上流テーブルのスキーマと異なる場合、上流テーブルに存在するが下流テーブルに欠落している nullable 列は自動的に既存のテーブルに追加されます。互換性のない列タイプや、欠落または冗長な NOT NULL 列など、自動的に適応できない場合、チェックは失敗し、エラーが報告されます。エラーメッセージで推奨される SQL ステートメントを実行してテーブルスキーマを変更できます。また、sink で sink.existing-table.schema-check.enabled を false に設定してチェックと自動列追加をスキップし、下流テーブルのスキーマに基づいてデータを書き込むこともできます。
例
以下の例は、典型的なシナリオでの設定を示しています。
Rest カタログへの書き込み
以下の例は、Paimon カタログが rest タイプの場合に Data Lake Formation (DLF) にデータを書き込む方法を示しています:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: ${mysql.source.table}
server-id: 8601-8604
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
catalog.properties プレフィックスを持つパラメーターについては、Flink CDC カタログ設定パラメーターをご参照ください。
FileSystem カタログへの書き込み
filesystem Paimon カタログを使用して Object Storage Service (OSS) に書き込むための設定例:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: ${mysql.source.table}
server-id: 8601-8604
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: filesystem
catalog.properties.warehouse: oss://default/test
catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
catalog.properties.fs.oss.accessKeyId: xxxxxxxx
catalog.properties.fs.oss.accessKeySecret: xxxxxxxx
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
catalog.properties プレフィックスを持つパラメーターについては、Paimon FileSystem カタログの作成をご参照ください。
パーティションテーブルへの書き込み
上流テーブルの TIMESTAMP 型の create_time フィールドを DATE フィールドに変換し、それを Paimon テーブルのパーティションキーとして使用します。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
transform:
- source-table: test_db.test_source_table
projection: \*, DATE_FORMAT(CAST(create_time AS TIMESTAMP), 'yyyy-MM-dd') as partition_key
primary-keys: id, create_time, partition_key
partition-keys: partition_key
description: add partition key
pipeline:
name: MySQL to Paimon Pipeline
単一テーブルの同期
リアルタイム性能と安定性に対する要件が高いテーブルには、単一テーブル同期ジョブを設定することを推奨します。単一テーブル同期は、複数テーブル同期の複雑さを回避します。以下の例は設定を示しています:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
データベース全体の同期
データベース全体を同期することで、設定および実行する必要のあるジョブの数が減り、O&M コストが削減されます。以下の例は設定を示しています:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
データベース全体の同期とデータベースおよびテーブル名の置換
データベース名とテーブル名を一括で置換するには、データベース全体の同期ジョブに route セクションを追加します。以下の例は設定を示しています:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
route:
# MySQL の test_db データベース内のすべてのテーブルを Paimon の test_db2 データベースに同期し、テーブル名を保持します。
- source-table: test_db.\.*
sink-table: test_db2.<>
replace-symbol: <>
pipeline:
name: MySQL to Paimon Pipeline
シャーディングされたデータベースとテーブルのマージ
複数のシャーディングされたテーブルからデータを単一の下流テーブルに書き込むには、以下のテンプレートを使用します:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.user\.*
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
route:
# MySQL の test_db データベース内のすべてのシャーディングされたテーブルは、単一の Paimon test_db.user テーブルにマージされます。
- source-table: test_db.user\.*
sink-table: test_db.user
pipeline:
name: MySQL to Paimon Pipeline
ジョブ再起動時に既存のテーブルを追加する
同期に既存のテーブルを追加するには、scan.newly-added-table.enabled = true を設定してジョブを再起動します。
ジョブが最初に scan.binlog.newly-added-table.enabled を true に設定して新しいテーブルをキャプチャした場合、その後 scan.newly-added-table.enabled を true に設定してジョブを再起動して既存のテーブルをキャプチャしないでください。そうしないと、データが複数回配信されます。
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
# ジョブ再起動時に、tables パラメーターでキャプチャされた新しいテーブルをチェックし、スナップショットを取得します。
# 注:scan.startup.mode: initial と一緒に使用する必要があります。
scan.newly-added-table.enabled: true
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
データベース全体の同期で特定のテーブルを除外する
テストテーブルやジョブのフェールオーバーを繰り返し引き起こすテーブルなど、特定のテーブルをデータベース全体の同期ジョブから除外するには、以下の設定を使用します:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
# この正規表現に一致するテーブルは同期されません。
tables.exclude: test_db.table1
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
データベース全体を Lance ファイル形式の Paimon テーブルに同期する
Lance は、ベクターおよびマルチモーダルデータ用に設計されたデータストレージフォーマットです。file.format テーブル作成パラメーターを設定して、Lance ファイル形式を使用する Paimon テーブルにデータを書き込みます。以下の例は設定を示しています:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.\.*
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
# 出力ファイルのファイル形式を lance に指定します
table.properties.file.format: lance
pipeline:
name: MySQL to Paimon Pipeline
EVOLVE モードの有効化
EVOLVE モードは、下流テーブルのスキーマを上流テーブルのスキーマと厳密に一致させます。これには、フィールドの削除やテーブルの削除などの操作が含まれます。このモードでは、下流がすべてのスキーマ変更イベントを適用できない場合、ジョブはフェールオーバーし、介入なしでは回復できない可能性があります。以下の例はジョブ設定を示しています:
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: flink
password: ${secret_values.password}
tables: test_db.test_source_table
server-id: 5401-5499
#(任意) 増分フェーズ中に新しく追加されたテーブルからデータを同期します。
scan.binlog.newly-added-table.enabled: true
#(任意) テーブルとフィールドのコメントを同期します。
include-comments.enabled: true
#(任意) TaskManager での潜在的な Out Of Memory (OOM) 問題を防ぐために、有界でないチャンクの配布を優先します。
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(任意) データ読み取りを高速化するために、解析とフィルタリングを有効にします。
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: paimon
name: Paimon Sink
using.built-in-catalog: paimon_dlf_catalog
#(任意) コミットユーザーを指定します。競合を避けるために、ジョブごとに異なるコミットユーザーを使用することを推奨します。
commit.user: your_job_name
#(任意) 読み取りパフォーマンスを向上させるために、削除ベクターを有効にします。
table.properties.deletion-vectors.enabled: true
pipeline:
name: MySQL to Paimon Pipeline
schema.change.behavior: evolve