Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Jointures Flink SQL

Dernière mise à jour :Aug 09, 2026

Flink SQL prend en charge des jointures complexes sur des tables dynamiques via diverses sémantiques de requête et types de jointure. Évitez les produits cartésiens : ils ne sont pas pris en charge et entraînent l'échec des requêtes. L'ordre des jointures n'est pas optimisé par défaut. Pour de meilleures performances, listez les tables dans la clause FROM par ordre croissant de fréquence de mise à jour.

Vue d'ensemble des jointures

Type de jointure

Description

Contraintes

Jointures régulières

La jointure régulière est le type de jointure le plus polyvalent. Tout nouvel enregistrement ou toute modification apportée à l'un ou l'autre côté de la jointure est visible et affecte l'ensemble du résultat.

La syntaxe d'une jointure régulière est flexible et prend en charge les opérations INSERT, UPDATE et DELETE sur les tables d'entrée. Toutefois, une jointure régulière conserve les données d'entrée des deux côtés dans son état, ce qui peut entraîner une croissance indéfinie. Pour limiter la taille de l'état, vous pouvez définir une durée de vie (TTL), mais cela risque de compromettre la précision des résultats.

Jointures d'intervalle

Une jointure d'intervalle renvoie le produit cartésien des enregistrements qui satisfont à la fois à la condition de jointure et aux contraintes temporelles.

Nécessite au moins une condition d'égalité et une condition de délimitation temporelle s'appliquant aux deux côtés. La plage de temps peut être définie par une condition unique (telle que <, <=, >= ou >), une condition BETWEEN ou une condition d'égalité sur les attributs de temps du même type (temps de traitement ou temps d'événement) dans les deux tables.

Jointures temporelles

Une jointure temporelle associe un flux à une table versionnée. Elle utilise le temps d'événement ou le temps de traitement pour faire correspondre les données à la version correcte de la table à un instant donné.

Les deux tables doivent utiliser la même sémantique temporelle (temps de traitement ou temps d'événement). Soyez attentif au cycle de vie du résultat de la jointure ; la condition de jointure repose généralement sur un horodatage spécifique.

Jointures de recherche

Une jointure de recherche enrichit généralement une table avec des données interrogées depuis un système externe. Elle nécessite qu'une table possède un attribut de temps de traitement et que l'autre soit sauvegardée par une table de dimension.

Nécessite qu'une table possède un attribut de temps de traitement, tandis que l'autre doit utiliser un connecteur source de recherche. Une condition d'égalité entre les deux tables est également requise.

Jointures latérales

Une jointure latérale relie une table à la sortie d'une fonction tabulaire. Chaque ligne de la table de gauche est jointe à toutes les lignes produites par l'appel correspondant à la fonction tabulaire.

Nécessite une condition de jointure ON TRUE constante dans la clause, telle que JOIN LATERAL TABLE(table_func(order_id)) t(res) ON TRUE.

Jointures régulières

Les quatre types courants de jointures régulières sont :

  • INNER JOIN : Renvoie les enregistrements qui satisfont à la condition de jointure dans les deux tables (intersection).

  • LEFT JOIN : Renvoie tous les enregistrements de la table de gauche, même s'il n'y a pas d'enregistrements correspondants dans la table de droite (préserve la table de gauche).

  • RIGHT JOIN : Renvoie tous les enregistrements de la table de droite, même s'il n'y a pas d'enregistrements correspondants dans la table de gauche (préserve la table de droite).

  • FULL OUTER JOIN : Renvoie l'union des deux tables, y compris les enregistrements correspondants et non correspondants.

Diagramme de jointure régulière

image

Exemple de jointure régulière

  1. Connectez-vous à la console Realtime Compute for Apache Flink.

  2. Dans la liste des espaces de travail, localisez l'espace de travail cible et cliquez sur Console dans la colonne Actions.

  3. Dans le volet de navigation de gauche, cliquez sur Development > ETL.

  4. Cliquez sur l'icône image, puis sur New Blank Stream Draft. Saisissez un name, sélectionnez une engine version et cliquez sur Create.

    Cet exemple montre comment utiliser une jointure pour corréler des lignes entre des tables, en associant les noms de code des super-héros à leurs véritables identités.
    CREATE TEMPORARY TABLE NOC (
      agent_id STRING,
      codename STRING
    )
    WITH (
      'connector' = 'faker',   -- The Faker connector generates simulated data.
      'fields.agent_id.expression' = '#{regexify ''(1|2|3|4|5){1}''}',  -- Generates a random number from 1 to 5.
      'fields.codename.expression' = '#{superhero.name}',   -- A built-in Faker function that randomly generates a superhero name.
      'number-of-rows' = '10'   -- Specifies that 10 rows of data are generated.
    );
    
    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}',  -- A built-in Faker function that randomly generates a name.
      'number-of-rows' = '10'
    );
    
    SELECT
        name,
        codename
    FROM NOC
    INNER JOIN RealNames ON NOC.agent_id = RealNames.agent_id;  -- If the agent_id (1-5) in both tables is equal, the name and codename are returned.
  5. Dans le coin supérieur droit, cliquez sur Debug, sélectionnez un cluster de débogage, puis cliquez sur OK. Si vous ne disposez pas de cluster de session, consultez la section Créer un cluster de session.

    Une fois le débogage terminé, le résultat est une table comportant deux colonnes, name et codename, affichant les données de sortie de la jointure régulière.

Pour plus d'informations sur l'utilisation des jointures régulières, consultez la section Instructions de jointure régulière.

Jointures d'intervalle

Une jointure d'intervalle associe des enregistrements provenant de deux flux qui se situent dans un intervalle de temps spécifique. Ce type de jointure sert généralement à corréler des événements provenant de deux flux qui se produisent à proximité dans le temps.

Diagramme de jointure d'intervalle

image

Exemple de jointure d'intervalle

Cet exemple joint les commandes aux expéditions. Il filtre les résultats pour inclure uniquement les commandes expédiées dans les trois heures suivant leur passation.

Créez un brouillon ETL comme décrit dans la section Jointures régulières.

CREATE TEMPORARY TABLE orders (
  id INT,
  order_time AS TIMESTAMPADD(HOUR, CAST(FLOOR(RAND()*(1-5+1)+5)*(-1) AS INT), CURRENT_TIMESTAMP)    -- Randomly gets a time from 2, 3, or 4 hours ago based on the local time.
)
WITH (
  'connector' = 'datagen',     -- The datagen connector can periodically generate random data.
  'rows-per-second'='10',      -- The rate of random data generation, 10 rows/s.
  'fields.id.kind'='sequence', -- The sequence generator.
  'fields.id.start'='1',       -- The sequence starts at 1.
  'fields.id.end'='100'        -- The sequence ends at 100.
);

CREATE TEMPORARY TABLE shipments (
  order_id INT,
  shipment_time AS TIMESTAMPADD(HOUR, CAST(FLOOR(RAND()*(1-5+1))+1 AS INT), CURRENT_TIMESTAMP)   -- Randomly gets a time from 0, 1, or 2 hours ago based on the local time.
)
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    -- The time difference between order_time and 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;  -- Filters for orders placed within three hours before the shipment time.

Dans le coin supérieur droit, cliquez sur Debug puis sur OK. Les résultats du débogage sont les suivants.

image

Pour plus d'informations sur les jointures d'intervalle, consultez la section Instructions de jointure d'intervalle.

Jointures temporelles

Une jointure temporelle est utilisée avec des tables dynamiques, dont le contenu change au fil du temps. Elle associe un flux d'événements à la version de la table dynamique qui était valide au moment de chaque événement.

Diagramme de jointure temporelle

image

Exemple de jointure temporelle

Cet exemple illustre un scénario commercial où les taux de change varient selon les périodes, et les commandes doivent être calculées en utilisant le taux en vigueur au moment de la transaction.

(Facultatif) Générer des données simulées

  1. Créez un brouillon ETL comme décrit dans la section Jointures régulières.

  1. Utilisez le connecteur Faker pour générer des données simulées et écrivez-les dans Kafka en tant que table de taux de change dynamique à l'aide du connecteur Upsert 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. Dans le coin supérieur droit, cliquez sur Deploy pour déployer le job.

  2. Dans le volet de navigation de gauche, cliquez sur O&M > Deployments. Localisez le job cible, cliquez sur Start dans la colonne Actions, sélectionnez Initial Mode, puis cliquez sur Start.

Créez un nouveau brouillon ETL comme décrit dans la section Jointures régulières pour lire les données simulées afin de les déboguer.

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 
;

Dans le coin supérieur droit, cliquez sur Debug puis sur OK.

Le taux de change a changé deux fois, à 20:16:11 et 20:35:22. Une transaction a eu lieu à 20:35:14. À ce moment-là, le taux de change n'avait pas encore été mis à jour. Par conséquent, le calcul de cette transaction doit utiliser le taux de change entré en vigueur à 20:16:11.

Jointures de recherche

Une jointure de recherche enrichit un flux de données avec des données statiques ou évoluant lentement provenant d'un système externe. Par exemple, vous pouvez joindre un flux de commandes en temps réel avec des données de référence, telles que les informations sur les produits stockées dans une base de données relationnelle. Cette jointure nécessite qu'une table possède un attribut de temps de traitement et que l'autre soit connectée via un connecteur de recherche, tel que le connecteur MySQL.

Diagramme de jointure de recherche

image
Remarque
  • Vous devez inclure FOR SYSTEM_TIME AS OF PROCTIME() pour joindre chaque enregistrement de flux à l'instantané de la table de dimension à ce moment-là. Si les données de la table de dimension changent ultérieurement, les résultats de la jointure ne sont pas affectés.

  • La condition ON doit inclure une condition d'égalité sur un champ que la table de dimension peut utiliser pour les recherches.

Exemple de jointure de recherche

Cet exemple montre comment enrichir les données de commande avec des données statiques provenant d'un connecteur externe pour ajouter les noms des produits.

Créez un brouillon ETL comme décrit dans la section Jointures régulières.

CREATE TEMPORARY TABLE orders ( 
    order_id STRING,
    product_id INT,
    order_total INT
) WITH (
  'connector' = 'faker',   -- The Faker connector generates simulated data.
  'fields.order_id.expression' = '#{Internet.uuid}',   -- Generates a random UUID.
  'fields.product_id.expression' = '#{number.numberBetween ''1'',''5''}',   -- Generates a random number from 1 to 5.
  'fields.order_total.expression' = '#{number.numberBetween ''1000'',''5000''}',  -- Generates a random number from 1000 to 5000.
  'number-of-rows' = '10'  -- The number of data rows to generate.
);

-- Connect to static product data in 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          -- Adds the is_importance field. The value is 1 if the order total exceeds 3000, indicating an important order.
FROM orders o
JOIN products FOR SYSTEM_TIME AS OF PROCTIME() AS p -- The FOR SYSTEM_TIME AS OF PROCTIME() clause ensures that as the join operator processes rows from `orders`, each row from `orders` is joined with the `products` rows that match the join condition.
ON o.product_id = p.product_id;

La table products contient 5 enregistrements d'exemple : 1-cola, 2-orange, 3-milk, 4-tea, 5-coffee.

Dans le coin supérieur droit, cliquez sur Debug puis sur OK. Les résultats du débogage sont les suivants.

image

Pour plus d'informations sur les jointures de recherche, consultez la section Instructions JOIN pour les tables de dimension.

Jointures latérales

Une jointure latérale utilise une sous-requête dans la clause FROM qui est exécutée pour chaque ligne de la requête externe. Cela peut améliorer la flexibilité et les performances des requêtes en réduisant les analyses de table. Toutefois, cette opération peut entraîner une dégradation des performances si la requête interne est complexe ou traite une grande quantité de données.

image

Exemple de jointure latérale

Cet exemple agrège les enregistrements de vente pour trouver les trois premiers produits par quantité totale vendue.

Créez un brouillon ETL comme décrit dans la section Jointures régulières.

CREATE TEMPORARY TABLE sale (
  sale_id STRING,
  product_id INT,
  sale_num INT
)
WITH (
  'connector' = 'faker',    -- The Faker connector generates simulated data.
  'fields.sale_id.expression' = '#{Internet.uuid}',   -- Generates a random UUID.
  'fields.product_id.expression' = '#{regexify ''(1|2|3|4|5){1}''}',   -- Generates a random number from 1 to 5.
  'fields.sale_num.expression' = '#{number.numberBetween ''1'',''10''}',  -- Generates a random integer from 1 to 10.
  'number-of-rows' = '50'   -- Generates 50 rows of data.
);

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;

Dans le coin supérieur droit, cliquez sur Debug puis sur OK. Les résultats du débogage sont les suivants.

image

Documents connexes