Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Jointure d'intervalle

Dernière mise à jour :Aug 09, 2026

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