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

Realtime Compute for Apache Flink:Paimon コネクタ

最終更新日:Sep 19, 2026

このトピックでは、ストリーミングデータレイクハウス用の 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

テーブルのコネクタを指定します。

文字列

いいえ

なし

  • Paimon カタログで Paimon テーブルを作成する場合、このパラメーターは不要です。

  • 別のカタログで Paimon 一時テーブルを作成する場合、このパラメーターを paimon に設定する必要があります。

path

テーブルのストレージパス。

文字列

いいえ

なし

  • Paimon カタログで Paimon テーブルを作成する場合、このパラメーターは不要です。

  • 別のカタログで Paimon 一時テーブルを作成する場合、このパラメーターは HDFS または OSS 内のテーブルのストレージディレクトリを指定します。

auto-create

指定されたパスにテーブルファイルが存在しない場合に、自動的に作成するかどうかを指定します。

ブール値

いいえ

false

有効な値:

  • false (デフォルト):指定されたパスに Paimon テーブルファイルが存在しない場合、ジョブは失敗します。

  • true:指定されたパスが存在しない場合、Flink は自動的に Paimon テーブルファイルを作成します。

file.format

データファイルのフォーマット。

文字列

いいえ

parquet

有効な値:

  • orc

  • parquet

  • avro

  • lance (Realtime Compute for Apache Flink 11.6 以降でサポート)

bucket

パーティションごとのバケット数。

整数

いいえ

1

Paimon は bucket-key に基づいてデータをバケットに分散します。

説明

各バケットには 5 GB 未満のデータを含めることを推奨します。

bucket-key

バケットキーとして使用される列。

文字列

いいえ

なし

データをバケットに分散するために使用される列を指定します。

複数の列名はカンマ (,) で区切ります。例:'bucket-key' = 'order_id,cust_id' は、order_id と cust_id 列に基づいてデータを分散します。

説明
  • このパラメーターが指定されていない場合、Paimon はプライマリキーに基づいてデータを分散します。

  • テーブルにプライマリキーがない場合、Paimon はすべての列の値に基づいてデータを分散します。

changelog-producer

Changelog の生成メカニズム。

文字列

いいえ

none

Paimon は任意の入力ストリームに対して完全な Changelog を生成できます。これは、すべての update_after レコードに対応する update_before レコードがあることを意味します。これにより、下流の消費が簡素化されます。有効な値:

  • none (デフォルト):追加の Changelog は生成されません。下流のコンシューマーは引き続きストリーミングモードで Paimon テーブルを読み取ることができますが、Changelog は不完全です (対応する update_before レコードなしで update_after レコードのみが含まれます)。

  • input:入力ストリームを Changelog ファイルに書き込み、これが完全な Changelog として機能します。

  • full-compaction:各フルコンパクション中に完全な Changelog を生成します。

  • lookup:各スナップショットがコミットされる前に完全な Changelog を生成します。

Changelog プロデューサーの選択方法の詳細については、Changelog の生成をご参照ください。

full-compaction.delta-commits

2つの連続するフルコンパクション間のコミットの最大数。

整数

いいえ

なし

フルコンパクションをトリガーする前に許可されるスナップショットコミットの最大数を指定します。

lookup.cache-max-memory-size

Paimon ディメンションテーブルのメモリキャッシュサイズ。

文字列

いいえ

256 MB

このパラメーターは、ディメンションテーブルのルックアップと lookup Changelog プロデューサーの両方のキャッシュサイズを制御します。

merge-engine

同じプライマリキーを持つレコードをマージするメカニズム。

文字列

いいえ

deduplicate

有効な値:

  • deduplicate:最新のレコードのみを保持します。

  • partial-update:既存のレコードを最新のレコードの null でない値で更新します。他の列は変更されません。

  • aggregation:指定された集計関数を使用して事前集計を実行します。

マージエンジンの詳細な分析については、マージエンジンをご参照ください。

partial-update.ignore-delete

削除 (-D) メッセージを無視するかどうかを指定します。

ブール値

いいえ

false

有効な値:

  • true:削除メッセージを無視します。

  • false:削除メッセージを処理します。潜在的な IllegalStateException または IllegalArgumentException エラーを防ぐために、sequence.field などのパラメーターを使用して削除を処理する戦略を設定する必要があります。

説明
  • Realtime Compute for Apache Flink 8.0.6 以前では、このパラメーターは merge-engine = 'partial-update' の部分更新シナリオでのみ有効です。

  • Realtime Compute for Apache Flink 8.0.7 以降では、このパラメーターは部分更新以外のシナリオとも互換性があり、ignore-delete パラメーターと同じ機能を持ちます。ignore-delete を代わりに使用することを推奨します。

  • このパラメーターを有効にするかどうかは、ビジネス要件と削除メッセージが期待されるものかどうかに基づいて決定してください。削除メッセージがジョブの意図したセマンティクスと一致しない場合は、早期に失敗させる方が良い選択であることが多いです。

ignore-delete

削除 (-D) メッセージを無視するかどうかを指定します。

ブール値

いいえ

false

有効な値は partial-update.ignore-delete のものと同じです。

説明
  • このパラメーターは、Realtime Compute for Apache Flink 8.0.7 以降でのみサポートされています。

  • このパラメーターは partial-update.ignore-delete と同じ機能を持ちます。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.time-retained 条件のいずれかが満たされ、かつ snapshot.num-retained.min 条件も満たされた場合に、スナップショットの有効期限切れがトリガーされます。

snapshot.num-retained.min

保持する最近のスナップショットの最小数。

整数

いいえ

10

N/A

snapshot.time-retained

スナップショットの保持期間。

文字列

いいえ

1h

この条件または snapshot.num-retained.max 条件のいずれかが満たされ、かつ snapshot.num-retained.min 条件も満たされた場合に、スナップショットの有効期限切れがトリガーされます。

write-mode

Paimon テーブルの書き込みモード。

文字列

いいえ

change-log

有効な値:

  • change-log:Paimon テーブルは、プライマリキーに基づいて挿入、削除、更新操作をサポートします。

  • append-only:Paimon テーブルは挿入操作のみを受け付け、プライマリキーをサポートしません。このモードは change-log モードよりも効率的です。

書き込みモードの詳細については、書き込みモードをご参照ください。

scan.infer-parallelism

Paimon ソーステーブルの並列度を自動的に推論するかどうかを指定します。

ブール値

いいえ

true

有効な値:

  • true:バケット数に基づいて Paimon ソーステーブルの並列度を自動的に推論します。

  • false:Realtime Compute for Apache Flink で設定されたデフォルトの並列度を使用します。エキスパートモードが有効な場合、ジョブは代わりに設定された並列度を使用します。

scan.parallelism

Paimon ソーステーブルの並列度。

整数

いいえ

なし

説明

ジョブの [構成] > [リソース] タブで [リソースモード] が [モード] に設定されている場合、このパラメーターは無視されます。

sink.parallelism

Paimon sink テーブルの並列度。

整数

いいえ

なし

説明

ジョブの [構成] > [リソース] タブで [リソースモード] が [モード] に設定されている場合、このパラメーターは無視されます。

sink.clustering.by-columns

Paimon sink テーブルへの書き込みのためのクラスタリング列を指定します。

文字列

いいえ

なし

Paimon の append-only テーブル (プライマリキーのないテーブル) の場合、このパラメーターはバッチジョブでのクラスタ化された書き込みを有効にします。このプロセスは、指定された列でデータをグループ化することにより、クエリのパフォーマンスを向上させます。

複数の列名はカンマ (,) で区切ります。例:'col1,col2'。

クラスタリングの詳細については、Apache Paimon 公式ドキュメントをご参照ください。

sink.delete-strategy

システムがリトラクションメッセージ (-D/-U) を正しく処理することを保証するための検証戦略を指定します。

​​

Enum

いいえ

NONE

有効な値と、リトラクションメッセージを処理する際の sink 演算子の期待される動作:

  • NONE (デフォルト):検証を実行しません。

  • IGNORE_DELETE:sink 演算子は -U および -D メッセージを無視する必要があります。リトラクションは発生しません。

  • NON_PK_FIELD_TO_NULL:sink 演算子は -U メッセージを無視する必要があります。-D メッセージについては、プライマリキーを保持し、他のすべての非プライマリキーフィールドを null に設定します。

    これは主に、複数の sink が同じテーブルに書き込む際の部分更新に使用されます。

  • DELETE_ROW_ON_PK:sink 演算子は -U メッセージを無視する必要がありますが、-D メッセージについてはプライマリキーに対応する行を削除します。

  • CHANGELOG_STANDARD:sink 演算子は、-U と -D の両方のメッセージについて、プライマリキーに対応する行を削除する必要があります。

説明
  • このパラメーターは、Realtime Compute for Apache Flink 8.0.8 以降でのみサポートされています。

  • リトラクションに対する実際の sink の動作は、ignore-delete や merge-engine などの他のパラメーターによって決定されます。このパラメーターは、実際の動作が選択された戦略と一致するかどうかを検証するだけです。不一致がある場合、ジョブはエラーメッセージとともに失敗し、修正を提案します。

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.nprobe

検索中にプローブする IVF クラスターの数。

整数

いいえ

16

値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。

ivf.refine_factor

top_k × refine_factor の IVF 候補を取得し、Paimon テーブルに保存されている元のベクターを使用して再ランキングします。

Float

いいえ

無効

すべての IVF バリアントでデフォルトで無効になっています。レイテンシーよりも再現率が優先されるシナリオで、圧縮インデックス (ivf-pq や ivf-hnsw-sq など) に最適です。

hnsw.ef_search

検索中の HNSW 検索幅。

整数

いいえ

0

値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。0 はネイティブライブラリのデフォルトが使用されることを示します。

diskann.search.list_size

Lumina DiskANN の検索リストサイズ。

整数

いいえ

max(1.5 × top_k, 16)

値が高いほど一般的に再現率は向上しますが、レイテンシーは増加します。

diskann.search.beam_width

Lumina DiskANN の検索ビーム幅。

整数

いいえ

4

—

search.parallel_number

並列 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, 23.0, 10, NULL>

  • <1, NULL, NULL, 'This is a book'>

  • <1, 25.2, NULL, NULL>

最初の列がプライマリキーの場合、最終的にマージされたレコードは <1, 25.2, 10, 'This is a book'> になります。

説明
  • 部分更新の結果をストリームで読み取るには、changelog-producer パラメーターを lookup または full-compaction に設定する必要があります。

  • 部分更新エンジンは削除メッセージを処理できません。partial-update.ignore-delete パラメーターを true に設定して削除メッセージを無視できます。

Aggregation

一部のユースケースでは、レコードの集計値のみが必要な場合があります。集約エンジンは、指定した集計関数を使用して同じプライマリキーを持つレコードを結合します。各非プライマリキー列について、fields.<field-name>.aggregate-function オプションを使用して集計関数を指定する必要があります。そうしないと、列はデフォルトで last_non_null_value 集計関数を使用します。たとえば、次の Paimon テーブル定義を考えます。

CREATE TABLE MyTable (
  product_id BIGINT,
  price DOUBLE,
  sales BIGINT,
  PRIMARY KEY (product_id) NOT ENFORCED
) WITH (
  'merge-engine' = 'aggregation',
  'fields.price.aggregate-function' = 'max',
  'fields.sales.aggregate-function' = 'sum'
);

price 列は max 関数で集計され、sales 列は sum 関数で集計されます。2つの入力レコード <1, 23.0, 15> と <1, 30.2, 20> が与えられた場合、最終結果は <1, 30.2, 35> になります。サポートされている集計関数とそれに対応するデータ型は次のとおりです:

  • sum:DECIMAL、TINYINT、SMALLINT、INTEGER、BIGINT、FLOAT、DOUBLE をサポートします。

  • min および max:DECIMAL、TINYINT、SMALLINT、INTEGER、BIGINT、FLOAT、DOUBLE、DATE、TIME、TIMESTAMP、TIMESTAMP_LTZ をサポートします。

  • last_value および last_non_null_value:すべてのデータ型をサポートします。

  • listagg:STRING をサポートします。

  • bool_and および bool_or:BOOLEAN をサポートします。

説明
  • sum 関数のみがリトラクションと削除をサポートし、他の集計関数はサポートしません。特定の列がリトラクションと削除メッセージを無視する必要がある場合は、'fields.${field_name}.ignore-retract'='true' を設定できます。

  • 集約の結果をストリームで読み取るには、changelog-producer パラメーターを lookup または full-compaction に設定する必要があります。

Changelog プロデューサー

changelog-producer パラメーターを設定して、Paimon が任意の入力ストリームに対して完全な Changelog (すべての update_after レコードに対応する update_before レコードがある) を生成するように設定します。次の表に、利用可能な Changelog プロデューサーを説明します。詳細については、Apache Paimon 公式ドキュメントをご参照ください。

プロデューサー

説明

None

changelog-producer を none (デフォルト) に設定すると、下流の Paimon ソーステーブルは特定のプライマリキーのデータの最新の状態のみを参照します。この不完全な Changelog は、コンシューマーがデータの以前の状態を判断できず、削除されたか、または最新の状態が何かしかわからないため、正しい計算を実行することを困難にします。

たとえば、下流のコンシューマーが列の合計を計算する必要があり、最新の値が 5 であることしかわからない場合、合計をどのように更新するかを判断できません。以前の値が 4 であれば、合計は 1 増加し、以前の値が 6 であれば、合計は 1 減少します。update_before レコードに敏感なコンシューマーは none プロデューサーを使用すべきではありませんが、他の Changelog プロデューサーはパフォーマンスのオーバーヘッドを伴います。

説明

データベースなどの下流のコンシューマーが update_before データに敏感でない場合は、none プロデューサーを使用できます。特定の要件に基づいて Changelog プロデューサーを設定してください。

Input

changelog-producer を input に設定すると、sink テーブルは入力ストリームを Changelog ファイルに二重書き込みします。

このプロデューサーは、入力ストリーム自体が既に完全な Changelog である場合 (Change Data Capture (CDC) からのデータなど) にのみ使用してください。

Lookup

changelog-producer を lookup に設定すると、sink テーブルはディメンションテーブルのルックアップに似たポイントクエリメカニズムを使用して、現在のスナップショットがコミットされる前に完全な Changelog を生成します。このプロデューサーは、任意の入力ストリームから完全な Changelog を生成します。

full-compaction プロデューサーと比較して、lookup プロデューサーは Changelog の適時性が優れていますが、全体的なリソース消費は多くなります。

高いデータ鮮度 (例:分単位) を必要とするユースケースには、このオプションを使用してください。

Full Compaction

changelog-producer を full-compaction に設定すると、sink テーブルは各フルコンパクション中に完全な Changelog を生成します。このプロデューサーは、任意の入力ストリームから完全な Changelog を生成します。フルコンパクションの間隔は full-compaction.delta-commits パラメーターで指定されます。

lookup プロデューサーと比較して、full-compaction プロデューサーはレイテンシーが高いですが、既存のフルコンパクションプロセスを活用するため、追加の計算は不要です。これにより、全体的なリソース消費が低くなります。

データ鮮度の要件が低いユースケース (例:時間単位) には、このオプションを使用してください。

書き込みモード

Paimon テーブルは、以下の書き込みモードをサポートしています。

モード

説明

Change-log

change-log 書き込みモードは、Paimon テーブルのデフォルトです。このモードは、プライマリキーに基づいて挿入、削除、更新操作をサポートします。このモードでは、マージエンジンと Changelog プロデューサーも使用できます。

Append-only

append-only 書き込みモードは、データの挿入のみをサポートし、プライマリキーを使用しません。このモードは change-log モードよりも効率的で、中程度のデータ鮮度要件 (例:分単位の鮮度) のシナリオでメッセージキューの代替として使用できます。

append-only 書き込みモードの詳細については、Apache Paimon 公式ドキュメントをご参照ください。このモードを使用する際は、次の点にご注意ください:

  • 要件に応じて bucket-key パラメーターを設定してください。そうしないと、Paimon テーブルはすべての列の値に基づいてデータをバケット化し、計算効率が悪くなります。

  • append-only 書き込みモードは、ある程度データの出力順序を保証できます。具体的な出力順序は次のとおりです:

    1. 異なるパーティションからのレコードの場合:scan.plan-sort-partition パラメーターが設定されている場合、値の小さいパーティションからのレコードが先に出力されます。そうでない場合は、先に作成されたパーティションからのレコードが先に出力されます。

    2. 同じパーティションとバケットからのレコードの場合、先に書き込まれたレコードが先に出力されます。

    3. 同じパーティションだが異なるバケットからのレコードの場合、異なるバケットは異なる並列タスクによって処理されるため、出力順序は保証されません。

Flink CDC データインジェストのターゲット

Paimon テーブルは、単一テーブルまたはデータベース全体からのデータのリアルタイム同期をサポートしています。上流のスキーマ変更は、リアルタイムで Paimon テーブルに同期されます。詳細については、データインジェストセクションをご参照ください。

Variant 読み取りプルーニング

クエリによって参照される Variant フィールドのみを読み取り、I/O とメモリのオーバーヘッドを削減します。

有効化の方法

パラメーター

デフォルト値

説明

variant.read.pushdown.enabled

false

true に設定して Variant 読み取りプルーニングを有効にします。

前提条件

  • 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)

パラメーター

パラメーター

タイプ

説明

source_table

STRING リテラル

Blob 列を含む Paimon テーブル。フォーマットは database.table です。テーブルのカタログが関数が属するカタログと同じ場合、catalog.database.table も使用できます。値は空でないリテラルでなければなりません。テーブルフィールド、動的カラム、CASE 式はサポートされていません。

descriptor

BYTES

BlobDescriptor バイト。blob-as-descriptor = 'true' が有効になった後に Blob 列を読み取った結果です。

validity

INTERVAL

URL の有効期間。値は正の秒数でなければなりません。例:INTERVAL '5' MINUTE。

戻り値

タイプ

説明

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

max_pt():最新のパーティションのみをロードします。max_two_pt():最新の2つのパーティションのみをロードします。

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

なし

値は paimon でなければなりません。

name

sink の名前。

いいえ

STRING

なし

catalog.properties.metastore

Paimon カタログのタイプ。

いいえ

STRING

filesystem

有効な値:

  • filesystem (デフォルト)

  • rest (Data Lake Formation (DLF) のみをサポートし、DLF-Legacy はサポートしません)

catalog.properties.*

Paimon カタログを作成するためのパラメーター。

いいえ

STRING

なし

詳細については、Paimon カタログの管理をご参照ください。

table.properties.*

Paimon テーブルを作成するためのパラメーター。

いいえ

STRING

なし

詳細については、Paimon テーブルオプションをご参照ください。

table.properties.file.format

テーブル内のデータファイルのストレージフォーマット。

String

いいえ

parquet

有効な値:

  • orc

  • parquet

  • avro

  • lance (VVR 11.6 以降)

catalog.properties.warehouse

ファイルストレージのルートディレクトリ。

いいえ

STRING

なし

このパラメーターは、catalog.properties.metastore が filesystem に設定されている場合にのみ適用されます。

commit.user-prefix

データファイルをコミットするためのユーザー名プレフィックス。

いいえ

STRING

なし

説明

異なるジョブには異なるユーザー名を設定することを推奨します。これにより、コミットの競合を引き起こしたジョブを特定しやすくなります。

partition.key

パーティションテーブルのパーティションキー。

いいえ

STRING

なし

異なるテーブルは ; で区切られ、異なるフィールドは , で区切られ、テーブルとフィールドは : で区切られます。例:testdb.table1:id1,id2;testdb.table2:name。

sink.cross-partition-upsert.tables

プライマリキーにすべてのパーティションキーが含まれていない、クロスパーティションアップサートが必要なテーブルをリストします。

いいえ

STRING

なし

クロスパーティション更新があるテーブルに適用されます。

  • フォーマット:テーブル名をセミコロン ; で区切ります。

  • パフォーマンスの推奨事項:この操作はリソースを大量に消費します。これらのテーブルには個別のジョブを作成してください。

重要
  • 基準を満たすすべてのテーブルをリストする必要があります。テーブル名を省略すると、データの重複が発生します。

sink.commit.parallelism

Commit 演算子の並列度を指定します。

いいえ

INTEGER

なし

Commit 演算子がボトルネックになっている場合、このパラメーターを使用して並列度を上げ、パフォーマンスを向上させます。

このパラメーターは、Realtime Compute for Apache Flink 11.6 以降でのみサポートされています。

説明

このパラメーターを設定すると、演算子の並列度が変更されます。ステートフルなジョブを再起動する際には、ジョブが部分的な演算子状態を無視できるように AllowNonRestoredState を指定する必要があります。

sink.existing-table.schema-check.enabled

既存のテーブルスキーマの互換性をチェックし、事前に作成されたテーブルにデータが書き込まれるときにスキーマを自動的に適応させるかどうかを指定します。

ブール値

いいえ

true

次のリストは、有効な値と、事前に作成されたテーブルにデータが書き込まれるときの動作を説明しています:

  • true (デフォルト):既存のテーブルの列タイプの互換性、および欠落または冗長な NOT NULL 列をチェックし、上流テーブルに存在するが下流テーブルに欠落している nullable 列を自動的に追加します。自動的に適応できない場合、ジョブは失敗し、問題を修正するために実行できる SQL ステートメントを返します。

  • false:前述のチェックと自動列追加をスキップし、下流テーブルの既存のスキーマに基づいてデータを書き込みます。最初の CREATE TABLE 文が実行されるときに行われる自動列追加のみがスキップされます。その後のスキーマ変更は自動的に適用されます。

説明

このパラメーターは、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

よくある質問