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

Realtime Compute for Apache Flink:マテリアライズドテーブルの作成と使用

最終更新日:Aug 07, 2026

このトピックでは、マテリアライズドテーブルの作成、データバックフィルの実行、データ鮮度の変更、データリネージの表示方法について説明します。

制限事項

  • この機能は、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

はい

マテリアライズドテーブルのデータ鮮度は、ソーステーブルからのデータ更新で許容される最大遅延を定義します。

説明
  • アップストリームテーブルもマテリアライズドテーブルである場合、ダウンストリームテーブルのデータ鮮度は、アップストリームテーブルのデータ鮮度の正の整数倍でなければなりません。

  • データ鮮度の最大値は 1 日です。

AS <select_statement>

はい

マテリアライズドテーブルにデータを投入するクエリを定義します。アップストリームテーブルには、マテリアライズドテーブル、通常のテーブル、またはビューを指定できます。SELECT ステートメントは、すべての Flink SQL クエリをサポートします。

PRIMARY KEY

いいえ

テーブル内の各行を一意に識別する列で、指定は任意です。これらの列には null 値を含めることはできません。

PARTITIONED BY

いいえ

マテリアライズドテーブルをパーティション分割するために使用される列で、指定は任意です。

WITH オプション

いいえ

テーブルのプロパティと、パーティションキーの時刻形式パラメーターを定義します。

たとえば、WITH ('partition.fields.#.date-formatter' = 'yyyyMMdd') を使用して、パーティション分割列の時間形式パラメーターを設定できます。パラメーターの使用方法の詳細については、手順の例をご参照ください。

REFRESH_MODE

いいえ

マテリアライズドテーブルのリフレッシュモードを指定します。指定されたリフレッシュモードは、フレームワークがデータ鮮度に基づいて自動的に推測するモードよりも優先されます。これにより、特定のシナリオに対応できます。

  • CONTINUOUS:ストリーミングジョブがマテリアライズドテーブルを増分更新します。ダウンストリームデータは、即時またはチェックポイント完了後に可視化されます。

  • FULL:ワークフローが定期的にマテリアライズドテーブルのバッチ更新をトリガーします。このモードでは、エンジンが完全更新か増分更新かを自動的に判断します。詳細については、「マテリアライズドテーブルの増分更新」をご参照ください。データのリフレッシュサイクルは、データ鮮度の設定と一致します。デフォルトでは、データはテーブルレベルで上書きされます。パーティションキーが存在する場合、最新のパーティションのみをリフレッシュするか、すべてのパーティションをリフレッシュするかを選択できます。

操作手順

  1. Realtime Compute for Apache Flink コンソールにログインします。

  2. 対象のワークスペースの[アクション]列で、[コンソール]をクリックします。

  3. 左側のナビゲーションペインで、[カタログ] をクリックし、目的の Apache Paimon カタログをクリックします。

  4. 対象のデータベースをクリックし、[マテリアライズドテーブルの作成] をクリックします。

    プライマリキー 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 オプションは、マテリアライズドテーブルの時間パーティション形式を指定します。# プレースホルダーは、文字列型のパーティショニング列の名前を表します。この情報により、システムは定期更新中に更新するパーティションを特定できます。

  5. マテリアライズドテーブルの更新を開始または停止します。

    1. カタログの下にある対象のマテリアライズドテーブルをクリックします。

    2. 右上隅で、[開始] または [停止] をクリックします。

      説明

      進行中の更新を停止した場合、ジョブは現在の更新サイクルを完了してから停止します。

  6. マテリアライズドテーブルのジョブの詳細を表示します。

    [テーブルスキーマ] タブの [基本情報] セクションで、[最新ジョブ] または [ワークフロー] の横にあるジョブ ID をクリックすると、詳細を表示できます。

マテリアライズドテーブルのクエリの変更

制限事項

  • クエリを変更できるのは、VVR 11.1 以降で作成されたマテリアライズドテーブルのみです。

  • クエリを変更する際、列の追加 と 計算ロジックの変更 のみが可能です。既存の列の順序を変更したり、その定義を変更したりすることはできません。

    操作

    サポート

    説明

    新しい列の追加

    はい

    既存の列の順序を維持したまま、スキーマに新しい列を追加できます。

    既存の列の計算ロジックの変更 (名前や型は変更しない)

    はい

    計算ロジックは変更できますが、列名とデータ型は同じでなければなりません。

    既存の列の順序の変更

    いいえ

    列の順序は固定です。変更するには、マテリアライズドテーブルを削除して再作成する必要があります。

    既存の列の名前またはデータ型の変更

    いいえ

    マテリアライズドテーブルを削除して再作成する必要があります。

例

  1. [テーブルの編集] をクリックしてクエリを変更します。コード例は次のとおりです。

    ALTER MATERIALIZED TABLE `paimon`.`default`.`mt-orders`
        AS
        SELECT
          *,
          price * quantity AS total_price
        FROM orders
        WHERE price * quantity > 1000
    ;
  2. [プレビュー] をクリックすると、変更前後の比較が表示されます。

    [マテリアライズドテーブルの変更詳細] ダイアログボックスには、次の変更が表示されます。[SQL ステートメント] エリアでは、SQL ステートメントが元の SELECT * から SELECT *, orders.price * orders.quantity AS total_price に変更され、新しい WHERE 条件 orders.price * orders.quantity > 1000 が追加されます。また、[テーブル列] エリアでは、新しい total_price (DOUBLE) フィールドが緑色でハイライト表示されます。

  3. [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 コストを削減します。

  1. カタログの下にある対象のマテリアライズドテーブルをクリックします。

  2. [データ情報] タブで、データをバックフィルします。

    マテリアライズドテーブルを作成したときにパーティションキーを定義した場合はパーティションテーブルであり、それ以外の場合は非パーティション化テーブルです。

    パーティションテーブル

    [パーティション] セクションで、初めてデータをバックフィルする場合、または必要なパーティションが存在しない場合は、[トリガー更新] をクリックします。 パーティションが既に存在する場合は、特定のパーティションを選択し、「操作」列の [更新] をクリックできます。

    [データ情報] タブをクリックし、ページ下部の [データパーティション] エリアで操作を実行します。

    パラメーター

    • パーティションキー:テーブルのパーティションキーです。たとえば、20241201 と入力すると、パーティション ds=20241201 内のすべてのデータがバックフィルされます。

    • タスク名:データバックフィルタスクの名前です。

    • 更新範囲 (オプション):ダウンストリームのマテリアライズドテーブルに更新をカスケードするかどうかを指定します。現在のテーブルから始まり、データリネージ内のすべてのマテリアライズドテーブルが更新されます。サポートされるダウンストリームリネージの最大深度は 6 レベルです。

      説明
      • パーティションテーブルを更新する場合、ダウンストリームのマテリアライズドテーブルは、開始テーブルとまったく同じパーティションキーを持つ必要があります。そうでない場合、更新操作は失敗します。

      • リネージ内のいずれかのマテリアライズドテーブルで更新が失敗した場合、後続のすべてのダウンストリームノードも失敗します。

    • デプロイターゲット:キューまたはセッションクラスターを選択できます。デフォルトの選択は default-queue です。

    非パーティション化テーブル

    [データ詳細] セクションで、[更新] をクリックします。

    パラメーター

    • タスク名:データバックフィルタスクの名前です。

    • 更新範囲:このオプションは非パーティション化テーブルでは利用できません。

      説明
      • 更新中、ダウンストリームのデータは完全にリフレッシュされます。

      • リネージ内のいずれかのマテリアライズドテーブルで更新が失敗した場合、後続のすべてのダウンストリームノードも失敗します。

      • 開始テーブルがストリーミングジョブによって更新される非パーティション化テーブルである場合、カスケード更新はサポートされません。

    • デプロイターゲット:キューまたはセッションクラスターを選択できます。デフォルトの選択は default-queue です。

  3. スケジュールおよびバッチによるバックフィル。

    タスクオーケストレーションを使用して、マテリアライズドテーブルのワークフローを作成し、スケジュールに基づいてバックフィルジョブを実行できます。また、ワークフローのデータバックフィル機能を使用して、指定した時間範囲のデータを一括でバックフィルすることもできます。

データ鮮度の変更

  1. 対応するカタログで、[マテリアライズドテーブル] データベースをクリックして、目的の [マテリアライズドテーブル] をクリックします。

  2. 右上隅で、[鮮度の変更] をクリックします。

    • マテリアライズドテーブルにプライマリキーがない場合、その更新方法をストリーミングとバッチの間で切り替えることはできません。システムは、データ鮮度の値が 30 分未満の場合はストリーミングジョブを使用し、30 分以上の場合はバッチジョブを使用します。したがって、プライマリキーのないテーブルでは、この 30 分のしきい値を超えて鮮度を変更することは許可されていません。

    • アップストリームテーブルがマテリアライズドテーブルである場合、ダウンストリームテーブルのデータ鮮度がアップストリームテーブルのデータ鮮度の正の整数倍であることを確認してください。

    • データ鮮度の最大値は 1 日です。

データリネージの表示

左側のナビゲーションペインで、[運用保守] > [データリネージ] を選択して、マテリアライズドテーブルのデータリネージページに移動します。このページでは、すべてのマテリアライズドテーブル間のリネージ関係を表示できます。また、マテリアライズドテーブルで直接 [更新の開始/停止] や [鮮度の変更] などの操作を実行することもできます。[詳細] をクリックすると、対応するマテリアライズドテーブルの詳細ページに移動します。

マテリアライズドテーブルのノードをクリックすると、その詳細パネルが展開されます。パネルには、[データ鮮度] の値、[最終更新時刻]、[マテリアライズドテーブルのステータス] が表示され、[更新をトリガー] アクションが利用できます。

関連ドキュメント