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

Realtime Compute for Apache Flink:Flink SQL の結合

最終更新日:Jul 25, 2026

Flink SQL は、さまざまなクエリセマンティクスと結合タイプを使用して、動的テーブルでの複雑な結合をサポートします。デカルト積の作成はサポートされておらず、クエリの失敗原因となるため避けてください。デフォルトでは結合順序は最適化されません。パフォーマンスを向上させるには、FROM 句でテーブルを更新頻度の低いものから高いものの順にリストしてください。

結合の概要

結合タイプ

説明

制約

通常結合

通常結合は、最も汎用的な結合タイプです。結合のどちらかの側の新しいレコードや変更は可視化され、結合結果全体に影響します。

通常結合の構文は柔軟で、入力テーブルに対する INSERT、UPDATE、DELETE 操作をサポートします。ただし、通常結合は両側の入力データをその状態に保持し、そのサイズは無期限に増大する可能性があります。過剰な状態サイズを防ぐために、状態の生存時間 (TTL) を設定できますが、これにより結果の精度が損なわれる可能性があります。

区間結合

区間結合は、結合条件と時間制約の両方を満たすレコードのデカルト積を返します。

少なくとも 1 つの等価結合条件と、両側に適用される時間制約条件が必要です。時間範囲は、単一の条件 (<、<=、>=、または > など)、BETWEEN 条件、または両方のテーブルで同じタイプの時間属性 (処理時間またはイベント時間) に対する等価条件として定義できます。

テンポラル結合

テンポラル結合は、ストリームをバージョン付きテーブルに結合し、イベント時間または処理時間のいずれかを使用して、特定の時点でデータを正しいテーブルバージョンに照合します。

両方のテーブルが同じ時間セマンティクス (処理時間またはイベント時間) を使用する必要があります。結合結果のライフサイクルに注意してください。結合条件は通常、特定のタイムスタンプに基づいています。

ルックアップ結合

ルックアップ結合は、通常、外部システムからクエリされたデータでテーブルをエンリッチするために使用されます。この結合では、一方のテーブルが処理時間属性を持ち、もう一方のテーブルがディメンションテーブルによってバックアップされている必要があります。

一方のテーブルが処理時間属性を持ち、もう一方のテーブルがルックアップソースコネクタを使用する必要があります。また、2 つのテーブル間の等価結合条件も必須です。

ラテラル結合

ラテラル結合は、テーブルをテーブル値関数の出力と接続します。左テーブルの各行は、テーブル値関数の対応する呼び出しによって生成されたすべての行と結合されます。

句に定数の ON TRUE 結合条件が必要です。例:JOIN LATERAL TABLE(table_func(order_id)) t(res) ON TRUE。

通常結合

通常結合の一般的な 4 つのタイプは次のとおりです。

  • INNER JOIN:両方のテーブルで結合条件を満たすレコード (共通部分) を返します。

  • LEFT JOIN:右テーブルに一致するレコードがない場合でも、左テーブルのすべてのレコードを返します (左テーブルを保持します)。

  • RIGHT JOIN:左テーブルに一致するレコードがない場合でも、右テーブルのすべてのレコードを返します (右テーブルを保持します)。

  • FULL OUTER JOIN:一致するレコードと一致しないレコードの両方を含む、両方のテーブルの和集合 (UNION) を返します。

通常結合の図

image

通常結合の例

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

  2. ワークスペースリストで、対象のワークスペースを見つけ、Actions 列の Console をクリックします。

  3. 左側のナビゲーションウィンドウで、[開発] > [ETL] をクリックします。

  4. image アイコンをクリックして [新規ブランクストリームドラフト] をクリックし、[名前] を入力し、[エンジンバージョン] を選択して、[作成] をクリックします。

    この例では、結合を使用してテーブル間の行を関連付け、スーパーヒーローのコードネームと本名を結び付ける方法を示します。
    CREATE TEMPORARY TABLE NOC (
      agent_id STRING,
      codename STRING
    )
    WITH (
      'connector' = 'faker',   -- Faker コネクタはシミュレーションデータを生成します。
      'fields.agent_id.expression' = '#{regexify ''(1|2|3|4|5){1}''}',  -- 1 から 5 までの乱数を生成します。
      'fields.codename.expression' = '#{superhero.name}',   -- スーパーヒーローの名前をランダムに生成する組み込みの Faker 関数。
      'number-of-rows' = '10'   -- 10 行のデータが生成されることを指定します。
    );
    
    CREATE TEMPORARY TABLE RealNames (
      agent_id STRING,
      name     STRING
    )
    WITH (
      'connector' = 'faker',
      'fields.agent_id.expression' = '#{regexify ''(1|2|3|4|5){1}''}',
      'fields.name.expression' = '#{Name.full_name}',  -- 名前をランダムに生成する組み込みの Faker 関数。
      'number-of-rows' = '10'
    );
    
    SELECT
        name,
        codename
    FROM NOC
    INNER JOIN RealNames ON NOC.agent_id = RealNames.agent_id;  -- 両方のテーブルの agent_id (1-5) が等しい場合、名前とコードネームが返されます。
  5. 右上隅で[デバッグ]をクリックし、デバッグクラスターを選択して、[OK] をクリックします。セッションクラスターがない場合は、「セッションクラスターを作成する」をご参照ください。

    デバッグが完了すると、結果は name と codename の 2 つの列を持つテーブルになり、通常結合の出力データが表示されます。

通常結合の使用方法の詳細については、「通常結合文」をご参照ください。

区間結合

区間結合は、特定の時間間隔内に収まる 2 つのストリームからのレコードを結合します。このタイプの結合は、通常、時間的に近接して発生する 2 つのストリームからのイベントを関連付けるために使用されます。

区間結合の図

区間結合の例

この例では、注文と出荷を結合します。結果をフィルタリングして、注文が行われてから 3 時間以内に出荷された注文のみを含めます。

「通常結合」で説明されているように ETL ドラフトを作成します。

CREATE TEMPORARY TABLE orders (
  id INT,
  order_time AS TIMESTAMPADD(HOUR, CAST(FLOOR(RAND()*(1-5+1)+5)*(-1) AS INT), CURRENT_TIMESTAMP)    -- ローカル時間に基づいて、2、3、または 4 時間前の時刻をランダムに取得します。
)
WITH (
  'connector' = 'datagen',     -- datagen コネクタは定期的にランダムなデータを生成できます。
  'rows-per-second'='10',      -- ランダムデータ生成レート、10 行/秒。
  'fields.id.kind'='sequence', -- シーケンスジェネレーター。
  'fields.id.start'='1',       -- シーケンスは 1 から始まります。
  'fields.id.end'='100'        -- シーケンスは 100 で終わります。
);

CREATE TEMPORARY TABLE shipments (
  order_id INT,
  shipment_time AS TIMESTAMPADD(HOUR, CAST(FLOOR(RAND()*(1-5+1))+1 AS INT), CURRENT_TIMESTAMP)   -- ローカル時間に基づいて、0、1、または 2 時間前の時刻をランダムに取得します。
)
WITH (
  'connector' = 'datagen',
  'rows-per-second'='5',
  'fields.order_id.kind'='sequence',
  'fields.order_id.start'='1',
  'fields.order_id.end'='100'
);

SELECT
  o.id AS order_id,
  o.order_time,
  s.shipment_time,
  TIMESTAMPDIFF(HOUR,o.order_time,s.shipment_time) AS hour_diff    -- order_time と shipment_time の時間差。
FROM orders o
JOIN shipments s ON o.id = s.order_id   
WHERE 
    o.order_time BETWEEN s.shipment_time - INTERVAL '3' HOUR AND s.shipment_time;  -- 出荷時刻の 3 時間前までに発注された注文をフィルタリングします。

右上隅で、[デバッグ] をクリックし、次に [OK] をクリックします。デバッグ結果は次のとおりです。

image

区間結合の詳細については、「区間結合文」をご参照ください。

テンポラル結合

テンポラル結合は、内容が時間とともに変化する動的テーブルで使用されます。イベントのストリームを、各イベントの時点で有効だった動的テーブルのバージョンに結合します。

テンポラル結合の図

テンポラル結合の例

この例では、為替レートが異なる期間で変化し、注文は取引時に有効だったレートを使用して計算する必要があるビジネスシナリオを示します。

(任意) シミュレーションデータの生成

  1. 「通常結合」で説明されているように ETL ドラフトを作成します。

  1. Faker コネクタを使用してシミュレーションデータを生成し、Upsert Kafka コネクタを使用して動的な為替レートテーブルとして Kafka に書き込みます。

CREATE TEMPORARY TABLE currency_rates (
  `currency_code` STRING,
  `eur_rate` DECIMAL(6,4),
  `rate_time` TIMESTAMP(3),
  WATERMARK FOR `rate_time` AS rate_time - INTERVAL '15' SECOND,
  PRIMARY KEY (currency_code) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'currency_rates',
  'properties.bootstrap.servers' = '${secret_values.kafkahost}',
  'key.format' = 'raw',
  'value.format' = 'json'
);

CREATE TEMPORARY TABLE transactions (
  `id` STRING,
  `currency_code` STRING,
  `total` DECIMAL(10,2),
  `transaction_time` TIMESTAMP(3),
  WATERMARK FOR `transaction_time` AS transaction_time - INTERVAL '30' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'transactions',
  'properties.bootstrap.servers' = '${secret_values.kafkahost}',
  'key.format' = 'raw',
  'key.fields' = 'id',
  'value.format' = 'json'
);

CREATE TEMPORARY TABLE currency_rates_faker (
  `currency_code` STRING,
  `eur_rate` DECIMAL(6,4),
  `rate_time` TIMESTAMP(3)
)
WITH (
  'connector' = 'faker',
  'fields.currency_code.expression' = '#{Currency.code}',
  'fields.eur_rate.expression' = '#{Number.randomDouble ''4'',''0'',''10''}',
  'fields.rate_time.expression' = '#{date.past ''15'',''SECONDS''}',
  'rows-per-second' = '2'
);

CREATE TEMPORARY TABLE transactions_faker (
  `id` STRING,
  `currency_code` STRING,
  `total` DECIMAL(10,2),
  `transaction_time` TIMESTAMP(3)
)
WITH (
  'connector' = 'faker',
  'fields.id.expression' = '#{Internet.UUID}',
  'fields.currency_code.expression' = '#{Currency.code}',
  'fields.total.expression' = '#{Number.randomDouble ''2'',''10'',''1000''}',
  'fields.transaction_time.expression' = '#{date.past ''30'',''SECONDS''}',
  'rows-per-second' = '2'
);

BEGIN STATEMENT SET;

INSERT INTO currency_rates
SELECT * FROM currency_rates_faker;

INSERT INTO transactions
SELECT * FROM transactions_faker;

END;
  1. 右上隅の [デプロイ] をクリックして、ジョブをデプロイします。

  2. 左側のナビゲーションウィンドウで、[O&M] > [デプロイメント] をクリックします。対象のジョブを見つけ、[アクション] 列の [開始] をクリックし、[初期モード] を選択してから、[開始] をクリックします。

「通常結合」セクションで説明されているように新しい ETL ドラフトを作成し、デバッグのためにシミュレーションデータを読み取ります。

CREATE TEMPORARY TABLE currency_rates (
  `currency_code` STRING,
  `eur_rate` DECIMAL(6,4),
  `rate_time` TIMESTAMP(3),
  WATERMARK FOR `rate_time` AS rate_time - INTERVAL '15' SECOND,
  PRIMARY KEY (currency_code) NOT ENFORCED
) WITH (
  'connector' = 'upsert-kafka',
  'topic' = 'currency_rates',
  'properties.bootstrap.servers' = '${secret_values.kafkahost}',
  'properties.auto.offset.reset' = 'earliest',
  'properties.group.id' = 'currency_rates',
  'key.format' = 'raw',
  'value.format' = 'json'
);

CREATE TEMPORARY TABLE transactions (
  `id` STRING,
  `currency_code` STRING,
  `total` DECIMAL(10,2),
  `transaction_time` TIMESTAMP(3),
  WATERMARK FOR `transaction_time` AS transaction_time - INTERVAL '30' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'transactions',
  'properties.bootstrap.servers' = '${secret_values.kafkahost}',
  'properties.auto.offset.reset' = 'earliest',
  'properties.group.id' = 'transactions',
  'key.format' = 'raw',
  'key.fields' = 'id',
  'value.format' = 'json'
);

SELECT
  t.id,
  t.total * c.eur_rate AS total_eur,
  c.eur_rate,
  t.total,
  c.currency_code,
  c.rate_time,
  t.transaction_time
FROM transactions t
JOIN currency_rates FOR SYSTEM_TIME AS OF t.transaction_time AS c
ON t.currency_code = c.currency_code 
;

右上隅で[デバッグ]をクリックし、[OK] をクリックします。

為替レートは 20:16:11 と 20:35:22 の 2 回変更されました。トランザクションは 20:35:14 に発生しました。その時点では、為替レートはまだ更新されていませんでした。したがって、このトランザクションの計算には、20:16:11 に有効になった為替レートを使用する必要があります。

ルックアップ結合

ルックアップ結合は、データストリームを外部システムの静的または緩やかに変化するデータでエンリッチします。たとえば、リアルタイムの注文ストリームを、リレーショナルデータベースに保存されている製品情報などの参照データと結合できます。この結合では、一方のテーブルが処理時間属性を持ち、もう一方のテーブルが MySQL コネクタなどのルックアップコネクタを介して接続されている必要があります。

ルックアップ結合の図

説明
  • 各ストリームレコードをその時点でのディメンションテーブルのスナップショットと結合するには、FOR SYSTEM_TIME AS OF PROCTIME() を含める必要があります。後でディメンションテーブルのデータが変更されても、結合結果には影響しません。

  • ON 条件には、ディメンションテーブルがルックアップに使用できるフィールドに対する等価条件を含める必要があります。

ルックアップ結合の例

この例では、外部コネクタからの静的データで注文データをエンリッチして、製品名を追加する方法を示します。

「通常結合」で説明されているように ETL ドラフトを作成します。

CREATE TEMPORARY TABLE orders ( 
    order_id STRING,
    product_id INT,
    order_total INT
) WITH (
  'connector' = 'faker',   -- Faker コネクタはシミュレーションデータを生成します。
  'fields.order_id.expression' = '#{Internet.uuid}',   -- ランダムな UUID を生成します。
  'fields.product_id.expression' = '#{number.numberBetween ''1'',''5''}',   -- 1 から 5 までの乱数を生成します。
  'fields.order_total.expression' = '#{number.numberBetween ''1000'',''5000''}',  -- 1000 から 5000 までの乱数を生成します。
  'number-of-rows' = '10'  -- 生成するデータ行の数。
);

-- MySQL の静的製品データに接続します。
CREATE TEMPORARY TABLE products (
 product_id INT,
 product_name STRING
)
WITH(
  'connector' = 'mysql',
  'hostname' = '${secret_values.mysqlhost}',
  'port' = '3306',
  'username' = '${secret_values.username}',
  'password' = '${secret_values.password}',
  'database-name' = 'db2024',
  'table-name' = 'products'
);

SELECT 
    o.order_id,
    p.product_name,
    o.order_total,
  CASE 
    WHEN o.order_total > 3000 THEN 1    
    ELSE 0
  END AS is_importance          -- is_importance フィールドを追加します。注文合計が 3000 を超える場合、値は 1 となり、重要な注文を示します。
FROM orders o
JOIN products FOR SYSTEM_TIME AS OF PROCTIME() AS p -- FOR SYSTEM_TIME AS OF PROCTIME() 句により、結合演算子が `orders` の行を処理する際に、`orders` の各行が結合条件に一致する `products` の行と結合されることが保証されます。
ON o.product_id = p.product_id;

products テーブルには、1-cola、2-orange、3-milk、4-tea、5-coffee の 5 つのサンプルレコードが含まれています。

右上隅で[デバッグ]をクリックし、次に[OK]をクリックします。デバッグ結果は次のとおりです。

image

ルックアップ結合の詳細については、「ディメンションテーブルの JOIN 文」をご参照ください。

ラテラル結合

ラテラル結合は、FROM 句でサブクエリを使用し、外部クエリの各行に対して実行されます。これにより、テーブルスキャンを減らすことで、クエリの柔軟性とパフォーマンスを向上させることができます。ただし、内部クエリが複雑であるか、大量のデータを処理する場合、この操作はパフォーマンスの低下につながる可能性があります。

ラテラル結合の例

この例では、販売記録を集計して、総販売数量で上位 3 つの製品を見つけます。

「通常結合」で説明されているように ETL ドラフトを作成します。

CREATE TEMPORARY TABLE sale (
  sale_id STRING,
  product_id INT,
  sale_num INT
)
WITH (
  'connector' = 'faker',    -- Faker コネクタはシミュレーションデータを生成します。
  'fields.sale_id.expression' = '#{Internet.uuid}',   -- ランダムな UUID を生成します。
  'fields.product_id.expression' = '#{regexify ''(1|2|3|4|5){1}''}',   -- 1 から 5 までの乱数を生成します。
  'fields.sale_num.expression' = '#{number.numberBetween ''1'',''10''}',  -- 1 から 10 までのランダムな整数を生成します。
  'number-of-rows' = '50'   -- 50 行のデータを生成します。
);

CREATE TEMPORARY TABLE products (
 product_id INT,
 product_name STRING,
 PRIMARY KEY(product_id) NOT ENFORCED
)
WITH(
  'connector' = 'mysql',
  'hostname' = '${secret_values.mysqlhost}',
  'port' = '3306',
  'username' = '${secret_values.username}',
  'password' = '${secret_values.password}',
  'database-name' = 'db2024',
  'table-name' = 'products'
);

SELECT 
    p.product_name,
    s.total_sales
FROM  products p
LEFT JOIN LATERAL
    (SELECT SUM(sale_num) AS total_sales FROM sale WHERE sale.product_id = p.product_id) s ON TRUE
    ORDER BY total_sales DESC
    LIMIT 3;

右上隅で[デバッグ]をクリックし、次に[OK] をクリックします。デバッグ結果は次のとおりです。

image

関連ドキュメント

  • イベントのイベント時間ではなく、処理時間のみに関心がある場合は、処理時間テンポラル結合を使用できます。詳細については、「処理時間テンポラル結合文」をご参照ください。

  • 通常結合の使用方法の詳細については、「通常結合文」をご参照ください。

  • 区間結合の使用方法の詳細については、「区間結合文」をご参照ください。

  • ルックアップ結合の使用方法の詳細については、「ディメンションテーブルの JOIN 文」をご参照ください。