Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Jointure temporelle en temps de traitement

Dernière mise à jour :Aug 09, 2026

Utilisez une jointure temporelle en temps de traitement pour enrichir un flux de données avec le dernier snapshot d'une table de dimension. Lors de l'exécution de la requête, chaque ligne entrante est mise en correspondance avec l'état actuel de la table de dimension, c'est-à-dire l'état au moment où la ligne arrive dans Realtime Compute for Apache Flink, et non l'horodatage encodé dans l'événement lui-même.

Cette approche diffère d'une jointure temporelle en temps d'événement, qui corrèle les lignes sur la base de l'horodatage intégré à l'événement.

Prérequis

Avant de commencer, assurez-vous que :

  • Vous utilisez Realtime Compute for Apache Flink avec Ververica Runtime (VVR) version 8.0.10 ou ultérieure.

  • Une table de dimension MySQL a été créée et est accessible.

Notes d'utilisation

  • Définissez execution.checkpointing.interval-during-backlog = 0 pour désactiver les points de contrôle pendant la synchronisation complète. Realtime Compute for Apache Flink ne prend pas en charge les points de contrôle lors de la synchronisation complète. Ce paramètre n'affecte pas les points de contrôle durant la synchronisation incrémentielle des données.

  • Définissez table.optimizer.proctime-temporal-join-strategy = TEMPORAL_JOIN pour activer les jointures temporelles en temps de traitement.

Syntaxe

La jointure temporelle en temps de traitement utilise la même syntaxe qu'une jointure standard avec une table de dimension :

SELECT column-names
FROM table1 [AS <alias1>]
[LEFT] JOIN table2 FOR SYSTEM_TIME AS OF PROCTIME() [AS <alias2>]
ON table1.column-name1 = table2.key-name1;

La clause FOR SYSTEM_TIME AS OF PROCTIME() indique à Flink d'interroger la table de dimension selon le temps de traitement de chaque ligne entrante.

Exemple

Cet exemple joint une table de faits Kafka à une table de dimension MySQL pour rechercher des numéros de téléphone par nom.

Données de test

Table 1: kafka_input

id (bigint) name (varchar) age (bigint)
1 Lee 22
2 Harry 20
3 Liban 28

Table 2: phoneNumber

name (varchar) phoneNumber (bigint)
David 1390000111
Brooks 1390000222
Liban 1390000333
Lee 1390000444

Code de test

-- Enable processing-time temporal joins.
SET 'table.optimizer.proctime-temporal-join-strategy' = 'TEMPORAL_JOIN';
-- Disable checkpointing during full synchronization.
SET 'execution.checkpointing.interval-during-backlog' = '0';

-- Define the Kafka fact table.
-- proc_time is a computed column that returns the current processing time.
CREATE TEMPORARY TABLE kafka_input (
  id        BIGINT,
  name      VARCHAR,
  age       BIGINT,
  proc_time AS PROCTIME()
) WITH (
  'connector'                    = 'kafka',
  'topic'                        = '<yourTopic>',
  'properties.bootstrap.servers' = '<yourKafkaBrokers>',
  'properties.group.id'          = '<yourKafkaConsumerGroupId>',
  'format'                       = 'csv'
);

-- Define the MySQL dimension table.
CREATE TEMPORARY TABLE phoneNumber (
  name        VARCHAR,
  phoneNumber BIGINT,
  PRIMARY KEY (name) NOT ENFORCED
) WITH (
  'connector'     = 'mysql',
  'hostname'      = '<yourHostname>',
  'port'          = '3306',
  'username'      = '<yourUsername>',
  'password'      = '<yourPassword>',
  'database-name' = '<yourDatabaseName>',
  'table-name'    = '<yourTableName>'
);

-- Define the output table (blackhole discards results; replace with your target sink).
CREATE TEMPORARY TABLE result_infor (
  id          BIGINT,
  phoneNumber BIGINT,
  name        VARCHAR
) WITH (
  'connector' = 'blackhole'
);

-- Join the Kafka stream with the MySQL dimension table at processing time.
INSERT INTO result_infor
SELECT
  t.id,
  w.phoneNumber,
  t.name
FROM kafka_input AS t
JOIN phoneNumber FOR SYSTEM_TIME AS OF t.proc_time AS w
ON t.name = w.name;

Remplacez les espaces réservés suivants par les valeurs réelles :

Espace réservé Description
<yourTopic> Nom du topic Kafka
<yourKafkaBrokers> Adresses des serveurs bootstrap Kafka
<yourKafkaConsumerGroupId> ID du groupe de consommateurs Kafka
<yourHostname> Nom d'hôte ou adresse IP du serveur MySQL
<yourUsername> Nom d'utilisateur MySQL
<yourPassword> Mot de passe MySQL
<yourDatabaseName> Nom de la base de données MySQL
<yourTableName> Nom de la table MySQL

Résultats des tests

id (bigint) phoneNumber (bigint) name (varchar)
1 1390000444 Lee
3 1390000333 Liban

La ligne avec id = 2 (Harry) ne possède aucune entrée correspondante dans la table de dimension phoneNumber ; elle est donc exclue des résultats. Utilisez LEFT JOIN au lieu de JOIN pour inclure les lignes sans correspondance avec un numéro de téléphone null.