Flink SQL prend en charge les fonctions de fenêtre TUMBLE, HOP et SESSION, basées sur le temps d'événement ou le temps de traitement.
Fonctions de fenêtre
Flink SQL permet d'agréger des fenêtres infinies sans définition explicite dans les instructions SQL. Il gère également les agrégations sur une fenêtre spécifique. Par exemple, pour compter les utilisateurs ayant cliqué sur une URL au cours de la minute précédente, définissez une fenêtre collectant les données de clic de cette dernière minute afin de calculer le résultat.
Flink SQL prend en charge les agrégats de fenêtre et les agrégats OVER. Cette rubrique traite des agrégats de fenêtre. Ces derniers utilisent deux attributs temporels (temps d'événement et temps de traitement) et sont compatibles avec les fonctions de fenêtre TUMBLE, HOP et SESSION.
Les agrégations de fenêtre (TUMBLE, HOP et SESSION) combinées aux fonctions LAST_VALUE, FIRST_VALUE ou TopN peuvent produire des résultats inexacts en raison des mécanismes de déclenchement des fenêtres ou de latences.
Attributs temporels
Flink SQL prend en charge deux attributs temporels : le temps d'événement et le temps de traitement. Le comportement du fenêtrage varie selon l'attribut utilisé.
-
Temps d'événement : correspond généralement à l'horodatage intégré à un enregistrement.
Une fenêtre se ferme lorsque le watermark dépasse l'heure de fin de la fenêtre. La sortie n'est générée qu'à l'arrivée des données de fermeture. Pour une sous-tâche unique, le watermark augmente de manière monotone. En présence de plusieurs sous-tâches ou tables source, Flink utilise la valeur minimale du watermark.
ImportantSi des enregistrements désordonnés existent ou si une sous-tâche ou une partition ne contient aucune donnée, le watermark ne peut pas progresser et la fenêtre risque de ne pas se fermer. Pour résoudre ce problème, spécifiez un décalage de watermark pour les données désordonnées et assurez-vous que toutes les sous-tâches et partitions reçoivent un flux de données. Si une partition est inactive, ajoutez
table.exec.source.idle-timeout: 10sdans le champ Other Configuration de la section Parameters sous l'onglet Configuration de la page Deployments. Consultez la rubrique Configuration pour plus de détails sur les paramètres.Après le traitement des données via GROUP BY, des opérations JOIN sur deux flux de données ou des nœuds de fenêtre OVER, la propriété de watermark est perdue et le temps d'événement ne peut plus servir au fenêtrage.
-
Temps de traitement : heure de l'horloge système au moment où Flink traite un événement.
Flink génère le temps de traitement, qui n'existe pas dans vos données brutes. Vous devez donc définir explicitement une colonne de temps de traitement.
RemarqueLe temps de traitement dépendant de la vitesse d'arrivée des événements et de l'ordre de traitement, les résultats de relecture peuvent varier d'une exécution à l'autre.
Agrégations de fenêtre en cascade
Une fois une agrégation de fenêtre terminée, la colonne rowtime perd son attribut de temps d'événement. Utilisez des fonctions utilitaires telles que TUMBLE_ROWTIME, HOP_ROWTIME ou SESSION_ROWTIME pour récupérer max(rowtime) de la fenêtre et l'utiliser comme nouvelle valeur rowtime. La valeur retournée égale window_end - 1, est de type TIMESTAMP et conserve l'attribut rowtime. Par exemple, pour la fenêtre [00:00, 00:15), la valeur 00:14:59.999 est retournée.
L'exemple suivant applique en cascade une fenêtre tumbling d'une heure au-dessus d'une fenêtre tumbling d'une minute :
CREATE TEMPORARY TABLE user_clicks(
username varchar,
click_url varchar,
eventtime varchar,
ts AS TO_TIMESTAMP(eventtime),
WATERMARK FOR ts AS ts - INTERVAL '2' SECOND -- Define a watermark for rowtime.
) with (
'connector'='sls',
...
);
CREATE TEMPORARY TABLE tumble_output(
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) with (
'connector'='datahub' -- Simple Log Service allows you to export only VARCHAR-type DDL statements. Therefore, DataHub is used to store data.
...
);
CREATE TEMPORARY VIEW one_minute_window_output AS
SELECT
TUMBLE_ROWTIME(ts, INTERVAL '1' MINUTE) as rowtime, -- Use TUMBLE_ROWTIME as the aggregation time of the level-two window.
username,
COUNT(click_url) as cnt
FROM user_clicks
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE),username;
BEGIN statement set;
INSERT INTO tumble_output
SELECT
TUMBLE_START(rowtime, INTERVAL '1' HOUR),
TUMBLE_END(rowtime, INTERVAL '1' HOUR),
username,
SUM(cnt)
FROM one_minute_window_output
GROUP BY TUMBLE(rowtime, INTERVAL '1' HOUR), username;
END;
Résultats intermédiaires
Les données intermédiaires de fenêtre se composent d'un état à clés et de données de minuteur, stockés dans différents backends. Choisissez une combinaison adaptée aux caractéristiques de votre tâche :
|
Stockage de l'état à clés |
Stockage des minuteurs |
|
Mémoire |
|
|
Mémoire |
|
|
Mémoire |
|
|
Fichier |
Les minuteurs servent principalement à déclencher les fenêtres expirées. Stockez-les en mémoire pour obtenir les meilleures performances. Si les minuteurs sont nombreux ou si la mémoire est limitée, utilisez RocksDBStateBackend pour les stocker dans un fichier RocksDB.