このトピックでは、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 テーブルにデータを直接挿入または更新できます。
-
Paimon プライマリキーテーブルは、INSERT、UPDATE_BEFORE、UPDATE_AFTER、DELETE を含むすべてのタイプのメッセージを受け入れることができます。書き込み時に、データマージメカニズムに基づいて、同じプライマリキーを持つデータをマージします。
-
Paimon Append-only テーブル (非主キーテーブルとも呼ばれます) は、INSERT タイプのメッセージのみを受け入れます。
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;
関連ドキュメント
-
Paimon テーブルへのデータの書き込みと消費の際、SQL ヒントを使用してテーブルパラメーターを一時的に変更できます。詳細については、「Paimon テーブルの管理」をご参照ください。
-
Paimon プライマリキーテーブルと Paimon Append-only テーブルの基本的な特徴と機能の詳細については、「Paimon プライマリキーテーブルと Append-only テーブル」をご参照ください。
-
さまざまなシナリオにおける Paimon プライマリキーテーブルと Append-only テーブルの一般的な最適化の詳細については、「Paimon パフォーマンスの最適化」をご参照ください。
-
Paimon テーブルからのデータ消費は スナップショットファイル に依存しています。スナップショットの有効期限が短すぎる、または消費ジョブが非効率である場合、消費中のスナップショットファイルが有効期限切れで削除され、消費ジョブで
File xxx not found, Possible causesエラーが報告されることがあります。解決策については、「Paimon 読み取りジョブでの「File xxx not found, Possible causes」エラーの解決」をご参照ください。