La jointure d'intervalle fusionne deux flux de données selon une clé commune. Elle associe les éléments de chaque flux lorsque leurs horodatages se situent dans un intervalle de temps relatif spécifié. Une fois les deux flux joints, les colonnes d'horodatage des flux d'entrée sont conservées, ce qui permet un traitement ultérieur basé sur le temps d'événement.
Syntaxe
SELECT column-names
FROM table1 [AS <alias1>]
[INNER | LEFT | RIGHT | FULL] JOIN table2
ON table1.column-name1 = table2.key-name1 AND TIMEBOUND_EXPRESSION
Types de jointure pris en charge : INNER JOIN (valeur par défaut lorsque JOIN est utilisé seul), LEFT JOIN, RIGHT JOIN et FULL JOIN. Les jointures SEMI JOIN et ANTI JOIN ne sont pas prises en charge.
L'expression TIMEBOUND_EXPRESSION doit lier l'horodatage d'un flux à un intervalle fermé relatif à l'horodatage de l'autre. Les formes d'expression suivantes sont prises en charge :
| Forme | Exemple |
|---|---|
| Égalité | ltime = rtime |
| Plage | ltime >= rtime AND ltime < rtime + INTERVAL '10' MINUTE |
| BETWEEN | ltime BETWEEN rtime - INTERVAL '10' SECOND AND rtime + INTERVAL '5' SECOND |
Exemple : mise en correspondance des commandes expédiées sous 4 heures
Cet exemple joint un flux de commandes à un flux d'expéditions pour identifier les commandes expédiées dans les 4 heures suivant leur passation.
Données de test
Table des commandes :
| id | productName | orderTime |
|---|---|---|
| 1 | phone | 2024-04-01 10:00:00.0 |
| 2 | laptop | 2024-04-01 10:02:00.0 |
| 3 | watch | 2024-04-01 10:03:00.0 |
| 4 | tablet | 2024-04-01 10:05:00.0 |
Table des expéditions :
| shipId | orderId | status | shiptime |
|---|---|---|---|
| 0 | 1 | shipped | 2024-04-01 11:00:00.0 |
| 1 | 2 | delivered | 2024-04-01 17:00:00.0 |
| 2 | 3 | shipped | 2024-04-01 12:00:00.0 |
| 3 | 4 | shipped | 2024-04-01 11:30:00.0 |
Instructions SQL
Les deux tables source utilisent un connecteur Kafka avec un filigrane retardé de 2 secondes sur la colonne de temps d'événement. La condition de jointure o.ordertime BETWEEN s.shiptime - INTERVAL '4' HOUR AND s.shiptime sélectionne les commandes dont l'heure de commande se situe dans les 4 heures précédant l'heure d'expédition.
CREATE TEMPORARY TABLE Orders(
id BIGINT,
productName VARCHAR,
ordertime TIMESTAMP(3),
WATERMARK wk FOR ordertime as withOffset(ordertime, 2000) -- 2-second delayed watermark on ordertime
) WITH (
'connector' = 'kafka',
'topic' = '<yourTopic>',
'properties.bootstrap.servers' = '<yourBrokers>',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
CREATE TEMPORARY TABLE Shipments(
shipId BIGINT,
orderId BIGINT,
status VARCHAR,
shiptime TIMESTAMP(3),
WATERMARK wk FOR shiptime as withOffset(shiptime, 2000) -- 2-second delayed watermark on shiptime
) WITH (
'connector' = 'kafka',
'topic' = '<yourTopic>',
'properties.bootstrap.servers' = '<yourBrokers>',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
-- MySQL sink
CREATE TEMPORARY TABLE rds_output(
id BIGINT,
productName VARCHAR,
status VARCHAR
) WITH (
'connector' = 'mysql',
'hostname' = '<yourHostname>',
'port' = '3306',
'username' = '<yourUsername>',
'password' = '<yourPassword>',
'database-name' = '<yourDatabaseName>',
'table-name' = '<yourTableName>'
);
INSERT INTO rds_output
SELECT id, productName, status
FROM Orders AS o
JOIN Shipments AS s ON o.id = s.orderId AND
o.ordertime BETWEEN s.shiptime - INTERVAL '4' HOUR AND s.shiptime;
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description |
|---|---|
<yourTopic> |
Nom du topic Kafka |
<yourBrokers> |
Adresses des serveurs bootstrap Kafka |
<yourHostname> |
Nom d'hôte MySQL |
<yourUsername> |
Nom d'utilisateur MySQL |
<yourPassword> |
Mot de passe MySQL |
<yourDatabaseName> |
Nom de la base de données cible |
<yourTableName> |
Nom de la table cible |
Résultat
| id (BIGINT) | productName (VARCHAR) | status (VARCHAR) |
|---|---|---|
| 1 | phone | shipped |
| 3 | watch | shipped |
| 4 | tablet | shipped |
La commande 2 (laptop) est exclue car son heure d'expédition (17:00) est postérieure de plus de 4 heures à l'heure de la commande (10:02).