このトピックでは、マテリアライズドテーブルの作成、データバックフィルの実行、データ鮮度の変更、データリネージの表示方法について説明します。
制限事項
-
この機能は、Ververica Runtime (VVR) 8.0.10 以降のバージョンでのみ利用できます。
-
マテリアライズドテーブルは、メタデータストレージとして Filesystem またはDLF を使用する Apache Paimon カタログでのみ作成できます。カスタムの Apache Paimon カタログはサポートされていません。
-
ジョブを開発およびデプロイする権限が必要です。詳細については、「開発コンソールの権限承認」をご参照ください。
-
一時テーブル、一時ユーザー定義関数、一時ビューなど、一時オブジェクトはサポートされていません。
マテリアライズドテーブルの作成
構文
CREATE MATERIALIZED TABLE [catalog_name.][db_name.]table_name
-- プライマリキー制約
[([CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED)]
[COMMENT table_comment]
-- パーティションキー
[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
-- With オプション
[WITH (key1=val1, key2=val2, ...)]
-- データ鮮度
FRESHNESS = INTERVAL '<num>' { SECOND | MINUTE | HOUR | DAY }
-- リフレッシュモード
[REFRESH_MODE = { CONTINUOUS | FULL }]
AS <select_statement>
パラメーター
|
パラメーター |
必須 |
説明 |
|
FRESHNESS |
はい |
マテリアライズドテーブルのデータ鮮度は、ソーステーブルからのデータ更新で許容される最大遅延を定義します。 説明
|
|
AS <select_statement> |
はい |
マテリアライズドテーブルにデータを投入するクエリを定義します。アップストリームテーブルには、マテリアライズドテーブル、通常のテーブル、またはビューを指定できます。 |
|
PRIMARY KEY |
いいえ |
テーブル内の各行を一意に識別する列で、指定は任意です。これらの列には null 値を含めることはできません。 |
|
PARTITIONED BY |
いいえ |
マテリアライズドテーブルをパーティション分割するために使用される列で、指定は任意です。 |
|
WITH オプション |
いいえ |
テーブルのプロパティと、パーティションキーの時刻形式パラメーターを定義します。 たとえば、 |
|
REFRESH_MODE |
いいえ |
マテリアライズドテーブルのリフレッシュモードを指定します。指定されたリフレッシュモードは、フレームワークがデータ鮮度に基づいて自動的に推測するモードよりも優先されます。これにより、特定のシナリオに対応できます。
|
操作手順
-
対象のワークスペースの[アクション]列で、[コンソール]をクリックします。
-
左側のナビゲーションペインで、[カタログ] をクリックし、目的の Apache Paimon カタログをクリックします。
-
対象のデータベースをクリックし、[マテリアライズドテーブルの作成] をクリックします。
プライマリキー
order_id、カテゴリ名order_name、日付フィールドdsを持つordersという名前のソーステーブルがあるとします。以下の例では、このテーブルに基づいてマテリアライズドテーブルを作成する方法を示します。-
ordersテーブルに基づいてマテリアライズドテーブルmt_orderを作成します。クエリはすべての列を選択し、データ鮮度は 5 秒に設定されます。CREATE MATERIALIZED TABLE mt_order FRESHNESS = INTERVAL '5' SECOND AS SELECT * FROM `paimon`.`db`.`orders` ; -
マテリアライズドテーブル
mt_orderに基づいて、マテリアライズドテーブルmt_idを作成します。クエリはorder_idとdsをテーブル列として選択し、order_idをプライマリキー、dsをパーティションキーとして設定し、データ鮮度を 30 分に設定します。CREATE MATERIALIZED TABLE mt_id ( PRIMARY KEY (order_id) NOT ENFORCED ) PARTITIONED BY(ds) FRESHNESS = INTERVAL '30' MINUTE AS SELECT order_id,ds FROM mt_order ; -
マテリアライズドテーブル
mt_orderに基づいてマテリアライズドテーブルmt_dsを作成し、dsパーティショニング列にdate-formatter(時刻フォーマット) を指定します。スケジュールが実行されるたびに、スケジューリング時刻から鮮度を引いた値が、対応するdsのパーティション値に変換されます。たとえば、データ鮮度が 1 時間に設定され、スケジューリング時刻が2024-01-01 00:00:00の場合、ds の計算値は 20231231 になり、パーティションds = '20231231'のデータのみが更新されます。スケジューリング時刻が2024-01-01 01:00:00の場合、ds の計算値は 20240101 になり、パーティションds = '20240101'のデータが更新されます。CREATE MATERIALIZED TABLE mt_ds PARTITIONED BY(ds) WITH ( 'partition.fields.ds.date-formatter' = 'yyyyMMdd' ) FRESHNESS = INTERVAL '1' HOUR AS SELECT order_id,order_name,ds FROM mt_order ;説明-
partition.fields.#.date-formatterでは、# プレースホルダーは STRING 型の有効なパーティショニング列でなければなりません。 -
partition.fields.#.date-formatterオプションは、マテリアライズドテーブルの時間パーティション形式を指定します。#プレースホルダーは、文字列型のパーティショニング列の名前を表します。この情報により、システムは定期更新中に更新するパーティションを特定できます。
-
-
-
マテリアライズドテーブルの更新を開始または停止します。
-
カタログの下にある対象のマテリアライズドテーブルをクリックします。
-
右上隅で、[開始] または [停止] をクリックします。
説明進行中の更新を停止した場合、ジョブは現在の更新サイクルを完了してから停止します。
-
-
マテリアライズドテーブルのジョブの詳細を表示します。
[テーブルスキーマ] タブの [基本情報] セクションで、[最新ジョブ] または [ワークフロー] の横にあるジョブ ID をクリックすると、詳細を表示できます。
マテリアライズドテーブルのクエリの変更
制限事項
-
クエリを変更できるのは、VVR 11.1 以降で作成されたマテリアライズドテーブルのみです。
-
クエリを変更する際、列の追加 と 計算ロジックの変更 のみが可能です。既存の列の順序を変更したり、その定義を変更したりすることはできません。
操作
サポート
説明
新しい列の追加
はい
既存の列の順序を維持したまま、スキーマに新しい列を追加できます。
既存の列の計算ロジックの変更 (名前や型は変更しない)
はい
計算ロジックは変更できますが、列名とデータ型は同じでなければなりません。
既存の列の順序の変更
いいえ
列の順序は固定です。変更するには、マテリアライズドテーブルを削除して再作成する必要があります。
既存の列の名前またはデータ型の変更
いいえ
マテリアライズドテーブルを削除して再作成する必要があります。
例
-
[テーブルの編集] をクリックしてクエリを変更します。コード例は次のとおりです。
ALTER MATERIALIZED TABLE `paimon`.`default`.`mt-orders` AS SELECT *, price * quantity AS total_price FROM orders WHERE price * quantity > 1000 ; -
[プレビュー] をクリックすると、変更前後の比較が表示されます。
[マテリアライズドテーブルの変更詳細] ダイアログボックスには、次の変更が表示されます。[SQL ステートメント] エリアでは、SQL ステートメントが元の
SELECT *からSELECT *,に変更され、新しい WHERE 条件orders.price*orders.quantityAStotal_priceが追加されます。また、[テーブル列] エリアでは、新しいorders.price*orders.quantity> 1000total_price(DOUBLE) フィールドが緑色でハイライト表示されます。 -
[OK] をクリックすると、[テーブルスキーマ] タブで新しく追加された列と変更されたクエリロジックを表示できます。
カラムの追加は、通常、ダウンストリームのコンシューマーに影響を与えません。ただし、ダウンストリームジョブが上流のマテリアライズドテーブルからデータを消費する際に、動的パース (たとえば SELECT * や自動フィールドマッピング) に依存している場合、ジョブが失敗したり、データ形式の不一致エラーが報告されたりする可能性があります。動的パースを避け、固定のカラムリストを使用し、上流のスキーマが変更された場合は速やかにダウンストリームテーブルのスキーマを更新することを推奨します。
増分更新
制限事項
この機能は、VVR 8.0.11 以降のバージョンでのみ利用できます。
更新モード
マテリアライズドテーブルは、ストリーミング、完全バッチ、増分バッチの 3 つの更新モードをサポートします。
モードはデータ鮮度の設定によって決まります。データ鮮度が 30 分未満の場合はストリーミングモードが、30 分以上の場合はバッチモードが有効になります。バッチモードでは、エンジンは完全更新か増分更新かを自動的に決定します。増分更新は、最後の更新以降に変更されたデータのみを計算し、それをマテリアライズドテーブルにマージします。完全更新は、テーブルまたはパーティション全体のデータを計算し、マテリアライズドテーブルの既存のデータを上書きします。バッチモードでは、エンジンは増分更新を優先し、増分更新が不可能な場合にのみ完全更新にフォールバックします。
増分更新の条件
増分更新は、マテリアライズドテーブルが以下のすべての条件を満たす場合にのみ実行されます。
-
テーブル定義に
partition.fields.#.date-formatterパラメーターが設定されていません。 -
ソーステーブルにプライマリキーが定義されていない。
-
マテリアライズドテーブル定義のクエリが、以下の表で説明されているように増分更新をサポートしている。
SQL 句
サポート
SELECT
列選択およびスカラー関数式 (ユーザー定義関数を含む) でサポートされます。集計関数はサポートされていません。
FROM
テーブル名とサブクエリでサポートされます。
WITH
共通テーブル式 (CTE) でサポートされます。
WHERE
さまざまなスカラー関数式 (ユーザー定義関数を含む) を含むフィルター条件でサポートされます。
WHERE [NOT] EXISTS <subquery>やWHERE <column> [NOT] IN <subquery>のようなサブクエリはサポートされていません。UNION
UNION ALLのみがサポートされます。JOIN
-
INNER JOINはサポートされます。 -
LEFT/RIGHT/FULL [OUTER] JOINは、以下で説明する特定のLATERAL JOINおよびルックアップ結合の場合を除き、サポートされていません。 -
テーブル関数式 (ユーザー定義関数を含む) を伴う
[LEFT [OUTER]] JOIN LATERALはサポートされます。 -
ルックアップ結合の場合、
A [LEFT [OUTER]] JOIN B FOR SYSTEM_TIME AS OF PROCTIME()のみがサポートされます。
説明-
JOINキーワードを使用しない暗黙的な結合 (例:SELECT * FROM a, b WHERE a.id = b.id) はサポートされます。 -
INNER JOINの増分計算では、両方のソーステーブルから完全なデータを読み取ります。
GROUP BY
サポートされていません。
-
増分更新の例
例1:スカラー関数を使用して orders ソーステーブルのデータを処理します。
CREATE MATERIALIZED TABLE mt_shipped_orders (
PRIMARY KEY (order_id) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
SELECT
order_id,
COALESCE(customer_id, 'Unknown') AS customer_id,
CAST(order_amount AS DECIMAL(10, 2)) AS order_amount,
CASE
WHEN status = 'shipped' THEN 'Completed'
WHEN status = 'pending' THEN 'In Progress'
ELSE 'Unknown'
END AS order_status,
DATE_FORMAT(order_ts, 'yyyyMMdd') AS order_date,
UDSF_ProcessFunction(notes) AS notes
FROM
orders
WHERE
status = 'shipped';
例2:LATERAL JOIN と ルックアップ結合 を使用して、orders ソーステーブルのデータをエンリッチします。
CREATE MATERIALIZED TABLE mt_enriched_orders (
PRIMARY KEY (order_id, order_tag) NOT ENFORCED
)
FRESHNESS = INTERVAL '30' MINUTE
AS
WITH o AS (
SELECT
order_id,
product_id,
quantity,
proc_time,
e.tag AS order_tag
FROM
orders,
LATERAL TABLE(UDTF_StringSplitFunction(tags, ',')) AS e(tag))
SELECT
o.order_id,
o.product_id,
p.product_name,
p.category,
o.quantity,
p.price,
o.quantity * p.price AS total_amount,
order_tag
FROM o
LEFT JOIN
product_info FOR SYSTEM_TIME AS OF PROCTIME() AS p
ON
o.product_id = p.product_id;
データバックフィル
以前は、履歴データのストリーム処理結果を修正するには、別のバッチジョブを開発する必要がありました。マテリアライズドテーブルを使用すると、履歴データパーティションを直接バックフィルできます。このアプローチはバッチ処理とストリーミング処理を統合し、開発と O&M コストを削減します。
-
カタログの下にある対象のマテリアライズドテーブルをクリックします。
-
[データ情報] タブで、データをバックフィルします。
マテリアライズドテーブルを作成したときにパーティションキーを定義した場合はパーティションテーブルであり、それ以外の場合は非パーティション化テーブルです。
パーティションテーブル
[パーティション] セクションで、初めてデータをバックフィルする場合、または必要なパーティションが存在しない場合は、[トリガー更新] をクリックします。 パーティションが既に存在する場合は、特定のパーティションを選択し、「操作」列の [更新] をクリックできます。
[データ情報] タブをクリックし、ページ下部の [データパーティション] エリアで操作を実行します。
パラメーター
-
パーティションキー:テーブルのパーティションキーです。たとえば、
20241201と入力すると、パーティションds=20241201内のすべてのデータがバックフィルされます。 -
タスク名:データバックフィルタスクの名前です。
-
更新範囲 (オプション):ダウンストリームのマテリアライズドテーブルに更新をカスケードするかどうかを指定します。現在のテーブルから始まり、データリネージ内のすべてのマテリアライズドテーブルが更新されます。サポートされるダウンストリームリネージの最大深度は 6 レベルです。
説明-
パーティションテーブルを更新する場合、ダウンストリームのマテリアライズドテーブルは、開始テーブルとまったく同じパーティションキーを持つ必要があります。そうでない場合、更新操作は失敗します。
-
リネージ内のいずれかのマテリアライズドテーブルで更新が失敗した場合、後続のすべてのダウンストリームノードも失敗します。
-
-
デプロイターゲット:キューまたはセッションクラスターを選択できます。デフォルトの選択は
default-queueです。
非パーティション化テーブル
[データ詳細] セクションで、[更新] をクリックします。
パラメーター
-
タスク名:データバックフィルタスクの名前です。
-
更新範囲:このオプションは非パーティション化テーブルでは利用できません。
説明-
更新中、ダウンストリームのデータは完全にリフレッシュされます。
-
リネージ内のいずれかのマテリアライズドテーブルで更新が失敗した場合、後続のすべてのダウンストリームノードも失敗します。
-
開始テーブルがストリーミングジョブによって更新される非パーティション化テーブルである場合、カスケード更新はサポートされません。
-
-
デプロイターゲット:キューまたはセッションクラスターを選択できます。デフォルトの選択は
default-queueです。
-
-
スケジュールおよびバッチによるバックフィル。
タスクオーケストレーションを使用して、マテリアライズドテーブルのワークフローを作成し、スケジュールに基づいてバックフィルジョブを実行できます。また、ワークフローのデータバックフィル機能を使用して、指定した時間範囲のデータを一括でバックフィルすることもできます。
データ鮮度の変更
-
対応するカタログで、[マテリアライズドテーブル] データベースをクリックして、目的の [マテリアライズドテーブル] をクリックします。
-
右上隅で、[鮮度の変更] をクリックします。
-
マテリアライズドテーブルにプライマリキーがない場合、その更新方法をストリーミングとバッチの間で切り替えることはできません。システムは、データ鮮度の値が 30 分未満の場合はストリーミングジョブを使用し、30 分以上の場合はバッチジョブを使用します。したがって、プライマリキーのないテーブルでは、この 30 分のしきい値を超えて鮮度を変更することは許可されていません。
-
アップストリームテーブルがマテリアライズドテーブルである場合、ダウンストリームテーブルのデータ鮮度がアップストリームテーブルのデータ鮮度の正の整数倍であることを確認してください。
-
データ鮮度の最大値は 1 日です。
-
データリネージの表示
左側のナビゲーションペインで、 を選択して、マテリアライズドテーブルのデータリネージページに移動します。このページでは、すべてのマテリアライズドテーブル間のリネージ関係を表示できます。また、マテリアライズドテーブルで直接 [更新の開始/停止] や [鮮度の変更] などの操作を実行することもできます。[詳細] をクリックすると、対応するマテリアライズドテーブルの詳細ページに移動します。
マテリアライズドテーブルのノードをクリックすると、その詳細パネルが展開されます。パネルには、[データ鮮度] の値、[最終更新時刻]、[マテリアライズドテーブルのステータス] が表示され、[更新をトリガー] アクションが利用できます。
関連ドキュメント
-
マテリアライズドテーブルの概要については、「マテリアライズドテーブルの管理」をご参照ください。
-
Paimon とマテリアライズドテーブルに基づいてバッチ処理とストリーミング処理を統合したデータレイクハウスで分析パイプラインを構築する方法、およびリアルタイム更新のためにデータ鮮度を変更してバッチモードからストリームモードに切り替える方法については、「クイックスタート:マテリアライズドテーブルを使用してバッチ処理とストリーミング処理を統合したデータレイクハウスを構築する」をご参照ください。