Flink SQL prend en charge les jointures régulières (également appelées jointures de flux doubles) sur des tables dynamiques. Contrairement aux jointures par lots, les résultats sont mis à jour en continu à mesure que de nouvelles données arrivent dans l'un ou l'autre flux, ce qui garantit que la sortie finale est cohérente avec celle qu'une jointure par lot produirait sur les mêmes données.
Contexte
Les jointures régulières mettent en mémoire tampon tous les enregistrements historiques des deux flux d'entrée dans l'état Flink, afin que toute ligne entrante puisse être comparée à toutes les lignes précédemment vues de l'autre côté. Cela signifie que l'état croît indéfiniment tant que le job est en cours d'exécution. Pour les jobs de longue durée traitant des flux à haut volume, une croissance illimitée de l'état peut épuiser la mémoire et provoquer des pannes.
Pour limiter la taille de l'état, définissez un délai d'expiration (TTL) pour l'état à l'aide du paramètre table.exec.state.ttl ou de l'indicateur JOIN_STATE_TTL spécifique à chaque flux (consultez la section Hints). Notez que l'éviction de l'état peut affecter l'exactitude des résultats : les lignes dont le TTL a expiré ne correspondront plus aux nouvelles arrivées.
Si la gestion d'un état illimité pose problème, envisagez également d'utiliser des jointures par intervalle, des jointures de recherche ou des jointures par fenêtre comme alternatives.
Syntaxe
tableReference [, tableReference ]*
| tableExpression [ NATURAL | INNER ] [ { LEFT | RIGHT | FULL } [ OUTER ] ] JOIN tableExpression [ joinCondition ]
| tableExpression CROSS JOIN tableExpression
| tableExpression [ CROSS | OUTER ] APPLY tableExpression
joinCondition:
ON booleanExpression
| USING '(' column [, column ]* ')'
tableReference: le nom de la table.tableExpression: l'expression de table.joinCondition: la condition de jointure.
Indiquez en premier la table ayant la fréquence de mise à jour la plus faible et en dernier celle ayant la fréquence la plus élevée. Cela réduit la quantité d'état à comparer pour chaque ligne entrante.
Types de jointure
INNER JOIN
Renvoie uniquement les lignes pour lesquelles la condition de jointure est satisfaite des deux côtés.
SELECT id, productName, status
FROM orders o
JOIN shipments s
ON o.id = s.orderId;
LEFT JOIN
Renvoie toutes les lignes de la table de gauche. Les lignes sans correspondance du côté droit génèrent des valeurs NULL pour les colonnes de droite.
SELECT o.id, o.productName, s.status
FROM orders o
LEFT JOIN shipments s
ON o.id = s.orderId;
RIGHT JOIN
Renvoie toutes les lignes de la table de droite. Les lignes sans correspondance du côté gauche génèrent des valeurs NULL pour les colonnes de gauche.
SELECT o.id, o.productName, s.status
FROM orders o
RIGHT JOIN shipments s
ON o.id = s.orderId;
FULL OUTER JOIN
Renvoie toutes les lignes des deux tables. Les lignes sans correspondance d'un côté ou de l'autre génèrent des valeurs NULL pour les colonnes non appariées.
SELECT o.id, o.productName, s.status
FROM orders o
FULL OUTER JOIN shipments s
ON o.id = s.orderId;
CROSS JOIN
Renvoie le produit cartésien des deux tables : chaque ligne de gauche est associée à chaque ligne de droite.
SELECT *
FROM orders o
CROSS JOIN shipments s;
Hints
Dans Ververica Runtime (VVR) 8.0.1 et versions ultérieures, utilisez hints pour spécifier des valeurs de délai d'expiration (TTL) différentes pour les états des flux de gauche et de droite. Cela réduit la quantité de données d'état à maintenir.
L'indicateur JOIN_STATE_TTL s'applique uniquement aux jointures régulières. Il ne prend pas en charge les jointures de recherche, par intervalle ou par fenêtre. Les indicateurs pour les jointures régulières sont expérimentaux. La syntaxe peut changer dans les versions futures.
Syntaxe
-- VVR 8.0.1 and later
SELECT /*+ JOIN_STATE_TTL('tableReference1' = 'ttl1' [, 'tableReference2' = 'ttl2']*) */ ...
-- VVR 8.0.7 and later also support the Apache Flink syntax
SELECT /*+ STATE_TTL('tableReference1' = 'ttl1' [, 'tableReference2' = 'ttl2']*) */ ...
Définissez tableReference sur un nom de table, un nom de vue ou un alias. Si vous spécifiez un alias pour une table, utilisez cet alias dans l'indicateur.
Si vous spécifiez un TTL pour un seul flux, l'autre flux utilise le TTL au niveau du déploiement défini par le paramètre table.exec.state.ttl. La valeur par défaut est de 1,5 jour. Pour plus de détails, consultez la section Basic configurations.
Exemples
Tous les exemples utilisent JOIN_STATE_TTL pour définir des TTL indépendants pour chaque flux. VVR 8.0.7 et versions ultérieures acceptent également la syntaxe équivalente STATE_TTL indiquée dans les commentaires.
Utilisation d'un alias
SELECT /*+ JOIN_STATE_TTL('o' = '3d', 'p' = '1d') */
o.rowtime, o.productid, o.orderid, o.units, p.name, p.unitprice
FROM Orders AS o
JOIN Products AS p
ON o.productid = p.productid;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('o' = '3d', 'p' = '1d') */
o.rowtime, o.productid, o.orderid, o.units, p.name, p.unitprice
FROM Orders AS o
JOIN Products AS p
ON o.productid = p.productid;
Utilisation d'un nom de table
SELECT /*+ JOIN_STATE_TTL('Orders' = '3d', 'Products' = '1d') */ *
FROM Orders
JOIN Products
ON Orders.productid = Products.productid;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('Orders' = '3d', 'Products' = '1d') */ *
FROM Orders
JOIN Products
ON Orders.productid = Products.productid;
Utilisation d'un nom de vue
CREATE TEMPORARY VIEW v AS
SELECT id, ...
FROM (
SELECT id, ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ..) AS rn
FROM src1
WHERE ...
) tmp
WHERE rn = 1;
SELECT /*+ JOIN_STATE_TTL('v' = '1d', 'b' = '3d') */ v.* , b.*
FROM v
LEFT JOIN src2 AS b ON v.id = b.id;
-- VVR 8.0.7 and later
SELECT /*+ STATE_TTL('v' = '1d', 'b' = '3d') */ v.* , b.*
FROM v
LEFT JOIN src2 AS b ON v.id = b.id;
Exemples
Exemple 1 : Jointure des tables Orders et Shipments
Données de test
Tableau 1. Orders
| id | productName | ordertime |
|---|---|---|
| 1 | phone_a | 2025-05-01 10:00:00.0 |
| 2 | notebook_x | 2025-05-01 10:02:00.0 |
| 3 | phone_b | 2025-05-01 10:03:00.0 |
| 4 | pad_m | 2025-05-01 10:05:00.0 |
Tableau 2. Shipments
| shipId | orderId | status | shipTime |
|---|---|---|---|
| 101 | 1 | shipped | 2025-05-01 11:00:00.0 |
| 102 | 2 | delivered | 2025-05-01 17:00:00.0 |
| 103 | 3 | shipped | 2025-05-01 12:00:00.0 |
| 104 | 4 | shipped | 2025-05-01 11:30:00.0 |
Instruction de test
SELECT id, productName, status
FROM orders o
JOIN shipments s
ON o.id = s.orderId;
Résultats du test
| id | productName | status |
|---|---|---|
| 1 | phone_a | shipped |
| 2 | notebook_x | delivered |
| 3 | phone_b | shipped |
| 4 | pad_m | shipped |
Exemple 2 : Jointure des tables datahub_stream1 et datahub_stream2
Données de test
Tableau 3. datahub_stream1
| a (BIGINT) | b (BIGINT) | c (VARCHAR) |
|---|---|---|
| 0 | 10 | test11 |
| 1 | 10 | test21 |
Tableau 4. datahub_stream2
| a (BIGINT) | b (BIGINT) | c (VARCHAR) |
|---|---|---|
| 0 | 10 | test11 |
| 1 | 10 | test21 |
| 0 | 10 | test31 |
| 1 | 10 | test41 |
Instruction de test
SELECT s1.c, s2.c
FROM datahub_stream1 AS s1
JOIN datahub_stream2 AS s2
ON s1.a = s2.a
WHERE s1.a = 0;
Résultats du test
| s1.c (VARCHAR) | s2.c (VARCHAR) |
|---|---|
| test11 | test11 |
| test11 | test31 |