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 = 0pour 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_JOINpour 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.