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

Realtime Compute for Apache Flink:Paimon テーブルへのデータの書き込みと消費

最終更新日:Sep 19, 2026

このトピックでは、Realtime Compute for Apache Flink コンソールを使用して Paimon テーブルのデータを挿入、更新、上書き、または削除する方法について説明します。また、これらのテーブルからデータを消費し、消費オフセットを指定する方法についても説明します。

前提条件

まず、Paimon カタログと Paimon テーブルを作成する必要があります。詳細については、「Paimon カタログの管理」をご参照ください。

制限事項

Paimon テーブルは、Ververica Runtime (VVR) 8.0.5 以降でのみサポートされています。

Paimon テーブルへのデータ書き込み

Flink CDC データ投入によるデータとスキーマの変更の同期

詳細については、「Paimon データ投入コネクター」をご参照ください。

INSERT INTO によるデータの挿入または更新

INSERT INTO ステートメントを使用して、Paimon テーブルにデータを直接挿入または更新できます。

INSERT OVERWRITE によるデータの上書き

上書き操作では、既存のデータがクリアされ、新しいデータが書き込まれます。INSERT OVERWRITE ステートメントを使用して、Paimon テーブル全体または特定のパーティションを上書きできます。以下のコードは、その例を示しています。

説明
  • INSERT OVERWRITE ステートメントは、バッチジョブでのみサポートされています。

  • デフォルトでは、INSERT OVERWRITE 操作はチェンジログデータを生成しません。ダウンストリームのストリーミングジョブは、削除・インポートされたデータを消費できません。これらのデータを消費する必要がある場合は、「INSERT OVERWRITE ステートメントの結果のストリーミングと消費」をご参照ください。

  • パーティション化されていないテーブル my_table 全体を上書きします。

    INSERT OVERWRITE my_table SELECT ...;
  • my_table テーブルの dt=20240108,hh=06 パーティションを上書きします。

    INSERT OVERWRITE my_table PARTITION (`dt` = '20240108', `hh` = '06') SELECT ...;
  • my_table テーブルのパーティションを動的に上書きします。SELECT ステートメントの結果に現れるパーティションは上書きされ、他のパーティションは変更されません。

    INSERT OVERWRITE my_table SELECT ...;
  • パーティション化されたテーブル my_table 全体を上書きします。

    INSERT OVERWRITE my_table /*+ OPTIONS('dynamic-partition-overwrite' = 'false') */ SELECT ...;

DELETE を使用したデータの削除

DELETE ステートメントを使用して、Paimon プライマリキーテーブルからデータを削除できます。DELETE ステートメントは、データ探索でのみ実行できます。

-- my_table テーブルから currency = 'UNKNOWN' のすべてのデータを削除します。
DELETE FROM my_table WHERE currency = 'UNKNOWN';

削除メッセージのフィルタリング

Paimon プライマリキーテーブルを使用する場合、デフォルトでは、DELETE メッセージは対応するプライマリキーを持つデータを削除します。Paimon テーブルでこれらのメッセージが処理されないようにするには、SQL ヒントを使用して次のパラメーターを true に設定します。これにより、削除メッセージがフィルタリングされます。

パラメーター

説明

デフォルト

ignore-delete

削除メッセージをフィルタリングするかどうかを指定します。

ブール値

false

シンクの並列度の調整

SQL ヒントを使用して次のパラメーターを設定し、シンクオペレーターの並列度を調整できます。

パラメーター

説明

デフォルト

sink.parallelism

Paimon シンクオペレーターの並列度を設定します。

整数

なし

たとえば、次の SQL ステートメントは、Paimon シンクオペレーターの並列度を 10 に設定します。

INSERT INTO t /*+ OPTIONS('sink.parallelism' = '10') */ SELECT * FROM s;

Paimon テーブルからのデータ消費

ストリーミングジョブ

説明

ストリーミングジョブで Paimon プライマリキーテーブルを消費する場合、チェンジログプロデューサーを設定する必要があります。

デフォルトでは、ストリーミングジョブの Paimon ソースオペレーターは、ジョブの開始時にまずテーブルの全量データを、その後はその時点からの差分データを生成します。

指定されたオフセットからのデータ消費

次のいずれかの方法で、指定されたオフセットからデータを消費できます:

  • ジョブの起動時に Paimon テーブルの全データを読み取る必要がなく、後続の増分データのみを読み取る必要がある場合は、SQL ヒントを介して 'scan.mode' = 'latest' を設定できます。

    SELECT * FROM t /*+ OPTIONS('scan.mode' = 'latest') */;
  • 完全なデータではなく、特定の時点からの増分データのみを消費したい場合は、SQL ヒントを使用して scan.timestamp-millis パラメーターを設定できます。パラメーター値は、Unix エポック (1970-01-01 00:00:00 UTC) から指定された時点までのミリ秒数を表します。

    SELECT * FROM t /*+ OPTIONS('scan.timestamp-millis' = '1678883047356') */;
  • 特定の時刻以降に書き込まれたデータから消費を開始し、その後継続的に差分データを消費し続けるには、次のいずれかの方法を使用します。

    説明

    この消費方法は、指定された時刻以降に変更されたデータファイルを読み取ります。コンパクションのため、データファイルには指定された時刻より前に書き込まれた少量のデータが含まれる場合があります。必要に応じて、SQL ジョブに WHERE 句を追加してデータをフィルタリングできます。

    • SQL ヒントは設定しません。ジョブを開始するときに、[Specify source's start time] を選択できます。[Job Start] ダイアログボックスで [Stateless start-up] を選択し、[Specify source's start time] スイッチを有効にして、ターゲット時間を設定します。

    • SQL ヒントを使用してscan.file-creation-time-millis パラメーターを設定します。

      SELECT * FROM t /*+ OPTIONS('scan.file-creation-time-millis' = '1678883047356') */;
  • 全データを消費するのではなく、特定の スナップショット ファイルから開始する増分データのみを消費したい場合は、SQL ヒントを使用して scan.snapshot-id パラメーターを設定できます。このパラメーターの値は、指定された スナップショット ファイルの ID です。

    SELECT * FROM t /*+ OPTIONS('scan.snapshot-id' = '3') */;
  • 特定のスナップショットファイルの全データを消費し、続けて増分データを消費する場合は、SQL ヒントを使用して 'scan.mode' = 'from-snapshot-full' および scan.snapshot-id パラメーターを設定できます。scan.snapshot-id パラメーターの値は、指定されたスナップショットファイルの ID です。

    SELECT * FROM t /*+ OPTIONS('scan.mode' = 'from-snapshot-full', 'scan.snapshot-id' = '1') */;

コンシューマー ID の指定

コンシューマー ID は、Paimon テーブルの消費進捗を保存します。これは主に次のシナリオで使用されます:

  • コンシューマー ID を設定すると、対応する消費進捗が Paimon テーブルのメタデータファイルに保存されます。これにより、後でステートレスモードで開始した場合でも、ジョブは中断した時点から消費を再開できます。

  • コンシューマー ID を設定すると、まだ消費されていないスナップショットは有効期限が切れても削除されません。これにより、消費速度がスナップショットの有効期限切れに追いつかず、エラーが発生するのを防ぎます。

ストリーミングジョブで consumer-id パラメーターを設定することで、Paimon ソースオペレーターにコンシューマー ID を割り当てることができます。コンシューマー ID の値には、任意の文字列を使用できます。コンシューマー ID が初めて作成されるとき、その開始オフセットは 指定されたオフセットから Paimon テーブルをコンシュームする のルールに従って決定されます。その後、同じコンシューマー ID を引き続き使用することで、Paimon テーブルの消費を再開できます。

例えば、次の SQL 文は、Paimon ソースオペレーターに対して test-id という名前のコンシューマー ID を設定します。特定のコンシューマー ID のオフセットをリセットしたい場合は、'consumer.ignore-progress' = 'true' も設定できます。

SELECT * FROM t /*+ OPTIONS('consumer-id' = 'test-id') */;
説明

コンシューマー ID によって消費されていないスナップショットファイルは、期限切れになっても削除されません。使用されなくなったコンシューマー ID がクリーンアップされない場合、対応するスナップショットファイルと履歴データファイルは削除されることはなく、ストレージ容量を消費します。consumer.expiration-time テーブルパラメーターを設定して、指定された期間使用されていないコンシューマー ID をクリーンアップできます。たとえば、'consumer.expiration-time' = '3d' は、3 日間使用されていないコンシューマー ID がクリーンアップされることを示します。

INSERT OVERWRITE 結果のストリーミングと消費

デフォルトでは、INSERT OVERWRITE 操作はチェンジログデータを生成せず、削除およびインポートされたデータはダウンストリームのストリーミングジョブでコンシュームできません。このようなデータをコンシュームする必要がある場合は、SQL ヒントを使用してストリーミングコンシュームジョブで 'streaming-read-overwrite' = 'true' を設定できます。

SELECT * FROM t /*+ OPTIONS('streaming-read-overwrite' = 'true') */;

バッチジョブ

デフォルトでは、バッチジョブの Paimon ソースオペレーターは最新のスナップショットを読み取り、Paimon テーブルの最新の状態データを出力します。

バッチタイムトラベル

SQL ヒントで scan.timestamp-millis パラメーターを設定することで、特定の時点における Paimon テーブルの状態を照会できます。パラメーター値は、Unix エポック (1970-01-01 00:00:00 UTC) から指定された時刻までのミリ秒数を表します。

SELECT * FROM t /*+ OPTIONS('scan.timestamp-millis' = '1678883047356') */;

SQL ヒントを使用して scan.snapshot-id パラメーターを指定されたスナップショットの ID に設定することで、スナップショットが作成された時点の Paimon テーブルの状態をクエリできます。

SELECT * FROM t /*+ OPTIONS('scan.snapshot-id' = '3') */;

スナップショット間の変更のクエリ

Paimon テーブル内の 2 つのスナップショット間でのデータ変更をクエリする場合、SQL ヒントを使用して incremental-between パラメーターを設定できます。 たとえば、スナップショット 20 とスナップショット 12 の間で変更されたすべてのデータを表示する場合、SQL 文は次のとおりです。

SELECT * FROM t /*+ OPTIONS('incremental-between' = '12,20') */;
説明

バッチジョブでは DELETE メッセージの消費がサポートされていないため、これらのメッセージはデフォルトで破棄されます。バッチジョブで DELETE メッセージを消費する場合は、監査ログシステムテーブルをクエリしてください。例: SELECT * FROM `t$audit_log ` /*+ OPTIONS('incremental-between' = '12,20') */;

ソースの並列度の調整

デフォルトでは、Paimon はパーティションやバケットの数などの情報に基づいてソースオペレーターの並列度を自動的に推論します。SQL ヒントを使用して次のパラメーターを設定し、並列度を調整できます。

パラメーター

デフォルト

説明

scan.parallelism

整数

なし

Paimon ソースオペレーターの並列度を設定します。

scan.infer-parallelism

ブール値

true

Paimon ソースオペレーターの並列度を自動的に推論するかどうかを指定します。

scan.infer-parallelism.max

整数

1024

自動的に推論される Paimon ソースオペレーターの並列度の上限。

次の SQL ステートメントは、Paimon ソースオペレーターの並列度を 10 に設定する例を示しています。

SELECT * FROM t /*+ OPTIONS('scan.parallelism' = '10') */;

Paimon テーブルのディメンションテーブルとしての使用

Paimon テーブルはディメンションテーブルとしても使用できます。ディメンションテーブルの JOIN 構文の詳細については、「ディメンションテーブルの JOIN ステートメント」をご参照ください。

VARIANT 型の書き込みと消費

Ververica Runtime (VVR) 11.1 以降では、Paimon テーブルは VARIANT 半構造化データ型をサポートします。この型を使用すると、PARSE_JSON または TRY_PARSE_JSON を使用して、VARCHAR JSON 文字列を VARIANT 型に変換できます。VARIANT 型を直接書き込んで消費することで、JSON のクエリと処理のパフォーマンスが大幅に向上します。

以下のコードは、その例を示しています:

CREATE TABLE `my-catalog`.`my_db`.`my_tbl` (
  k BIGINT,
  info VARIANT
);
INSERT INTO `my-catalog`.`my_db`.`my_tbl` 
SELECT k, PARSE_JSON(jsonStr) FROM T;

関連ドキュメント