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 |
|
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. |
|
|
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 |
|
|
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. |
|
|
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. |
|
|
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 |
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

Exemple de jointure régulière
Connectez-vous à la console Realtime Compute for Apache Flink.
Dans la liste des espaces de travail, localisez l'espace de travail cible et cliquez sur Console dans la colonne Actions.
Dans le volet de navigation de gauche, cliquez sur .
-
Cliquez sur l'icône
, 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. -
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
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.

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
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.
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
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.

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.
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.

Documents connexes
Si vous vous souciez uniquement du temps de traitement des événements et non de leur temps d'événement, vous pouvez utiliser une jointure temporelle de temps de traitement. Pour plus d'informations, consultez la section Instructions de jointure temporelle de temps de traitement.
Pour plus d'informations sur l'utilisation des jointures régulières, consultez la section Instructions de jointure régulière.
Pour plus d'informations sur l'utilisation des jointures d'intervalle, consultez la section Instructions de jointure d'intervalle.
Pour plus d'informations sur l'utilisation des jointures de recherche, consultez la section Instructions JOIN pour les tables de dimension.