Realtime Compute for Apache Flink prend en charge deux types d'agrégation par fenêtre : l'agrégation par fenêtre de groupe et l'agrégation par fonction table (TVF) de fenêtre. Cette rubrique décrit leur syntaxe, les scénarios dans lesquels l'agrégation TVF de fenêtre revient à un mode sans TVF, ainsi que la prise en charge des flux de mise à jour selon les types de fenêtre.
Choisir entre les deux syntaxes
| Agrégation par fenêtre de groupe | Agrégation TVF de fenêtre | |
|---|---|---|
| Opérateur | GroupWindowAggregation |
WindowAggregate |
| Fonctions de fenêtre | TUMBLE, HOP, SESSION | TUMBLE, HOP, CUMULATE, SESSION |
| Statut | Obsolète | Recommandée |
| Optimisations des performances | Non | Oui |
Prise en charge de GROUPING SETS |
Non | Oui |
| Top-N de fenêtre après agrégation | Non | Oui |
| Prise en charge des flux de mise à jour | Oui (VVR) | Oui (VVR, tous les types de fenêtre) |
Utilisez l'agrégation TVF de fenêtre. Elle prend en charge tous les types de fenêtre disponibles avec l'agrégation par fenêtre de groupe, ainsi que CUMULATE. Elle offre des optimisations de performances, la prise en charge de GROUPING SETS, et vous permet d'appliquer Window Top-N aux résultats d'agrégation.
Agrégation par fenêtre de groupe (obsolète)
L'agrégation par fenêtre de groupe définit les fenêtres dans la clause GROUP BY. Elle correspond à l'opérateur GroupWindowAggregation et prend en charge les fonctions de fenêtre TUMBLE, HOP et SESSION.
Pour connaître la syntaxe, des exemples et les détails des fonctionnalités, consultez Agrégation par fenêtre de groupe.
Changement de comportement dans VVR 11.x pour l'agrégation par fenêtre de groupe
À partir de VVR 11.x (Flink 1.20), le système ne réécrit plus automatiquement l'agrégation par fenêtre de groupe (syntaxe obsolète) en un plan d'exécution d'agrégation TVF de fenêtre.
Dans VVR 8.x, le système réécrivait automatiquement la syntaxe obsolète vers le nouveau plan d'exécution, ce qui activait l'optimisation d'agrégation en deux phases Local-Global. À partir de VVR 11.x, cette réécriture automatique n'est plus le comportement par défaut : la syntaxe obsolète conserve son plan d'exécution physique d'origine.
Impact : les jobs utilisant la syntaxe obsolète ne bénéficient plus automatiquement de l'optimisation d'agrégation en deux phases. Les performances peuvent se dégrader dans les scénarios impliquant de grands volumes de données ou une distribution inégale des données.
Migrez l'agrégation par fenêtre de groupe vers l'agrégation TVF de fenêtre (nouvelle syntaxe). Pour connaître les conditions d'activation de l'optimisation d'agrégation en deux phases, consultez Optimisation Local-Global pour l'agrégation par fenêtre.
Compatibilité avec le comportement précédent
Si la migration n'est pas immédiatement possible, activez le paramètre suivant pour restaurer le comportement de VVR 8.x :
| Paramètre | Description | Valeur par défaut |
|---|---|---|
table.optimizer.window-rewrite-enabled |
Active la réécriture automatique de la syntaxe obsolète vers la nouvelle syntaxe, ce qui active l'optimisation d'agrégation en deux phases. | false (depuis VVR 11.x) |
Exemple de configuration (à ajouter aux paramètres du job) :
table.optimizer.window-rewrite-enabled: true
Ce paramètre sert uniquement de mesure de compatibilité transitoire lors des mises à niveau. La syntaxe obsolète pourrait perdre la prise en charge de la réécriture dans les versions futures. Migrez vers la nouvelle syntaxe dès que possible.
Agrégation TVF de fenêtre
L'agrégation TVF de fenêtre définit les fenêtres via une clause GROUP BY incluant les colonnes window_start et window_end produites par les TVF de fenêtre. Elle correspond à l'opérateur WindowAggregate et prend en charge les fonctions de fenêtre TUMBLE, HOP, CUMULATE et SESSION.
Contrairement à l'agrégation sur des tables continues, l'agrégation TVF de fenêtre ne produit aucun résultat intermédiaire : seul un résultat final est émis à la fin de chaque fenêtre. Les données d'état intermédiaires sont nettoyées automatiquement.
Pour connaître la syntaxe, des exemples et les détails des fonctionnalités, consultez Agrégation TVF de fenêtre.
Syntaxe TVF de fenêtre SESSION : VVR 11.x vs VVR 8.x
La syntaxe TVF de fenêtre SESSION diffère selon les versions de VVR. Effectuez la mise à niveau vers VVR 11.1 ou ultérieure pour utiliser la syntaxe complète.
VVR 11.x (Flink 1.20)
SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap)
| Paramètre | Description |
|---|---|
data |
Une table comportant une colonne d'attribut temporel |
keycols |
(Facultatif) Colonnes utilisées pour partitionner les données avant le fenêtrage de session |
timecol |
Colonne d'attribut temporel mappée aux fenêtres de session |
gap |
Intervalle de temps maximal entre deux événements appartenant à la même session |
VVR 8.x (Flink 1.17)
SESSION(TABLE data, DESCRIPTOR(timecol), gap)
| Paramètre | Description |
|---|---|
data |
Une table comportant une colonne d'attribut temporel |
timecol |
Colonne d'attribut temporel mappée aux fenêtres de session |
gap |
Intervalle de temps maximal entre deux événements appartenant à la même session |
VVR 8.x ne prend pas en chargePARTITION BY. Les champs de partitionnement sont déduits implicitement à partir de la clauseGROUP BY.
Comparaison des syntaxes SESSION : VVR 11.x vs VVR 8.x
| VVR 11.x | VVR 8.x | |
|---|---|---|
| Syntaxe | SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap) |
SESSION(TABLE data, DESCRIPTOR(timecol), gap) |
| Spécification des champs de partitionnement | Explicite — via PARTITION BY(keycols) |
Implicite — via la clause GROUP BY |
| Restrictions sur les champs de partitionnement | Aucune | Doivent figurer dans GROUP BY ; ne peuvent pas être window_start, window_end ni window_time |
Utilisation autonome de SESSION() |
Prise en charge | Doit être utilisée avec GROUP BY |
| Fusion de la fonction de fenêtre avec l'agrégation | Prise en charge | Non prise en charge — l'agrégation doit correspondre aux champs de partitionnement |
Les exemples suivants sont équivalents. Tous deux utilisent item comme champ de partitionnement.
-- The Bid table schema (used in all examples below)
> desc Bid;
+-------------+------------------------+------+-----+--------+---------------------------------+
| name | type | null | key | extras | watermark |
+-------------+------------------------+------+-----+--------+---------------------------------+
| bidtime | TIMESTAMP(3) *ROWTIME* | true | | | `bidtime` - INTERVAL '1' SECOND |
| price | DECIMAL(10, 2) | true | | | |
| item | STRING | true | | | |
+-------------+------------------------+------+-----+--------+---------------------------------+
-- VVR 11.x: partition field declared explicitly in SESSION()
> SELECT window_start, window_end, item, SUM(price) AS total_price
FROM TABLE(
SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
GROUP BY item, window_start, window_end;
-- VVR 8.x: partition field inferred from GROUP BY
> SELECT window_start, window_end, item, SUM(price) AS total_price
FROM TABLE(
SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
GROUP BY item, window_start, window_end;
| VVR 11.x | VVR 8.x | |
|---|---|---|
| Partitionnement de la fenêtre SESSION | SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES) |
SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES) |
| Agrégation et fusion des fenêtres | Fusion directe prise en charge (par exemple, SUM(price) au sein de la fenêtre) |
Les champs d'agrégation doivent correspondre aux champs de partitionnement de la fenêtre (par exemple, GROUP BY item) |
Quand l'agrégation TVF de fenêtre revient à un mode sans TVF
Lorsqu'une requête inclut une TVF de fenêtre mais ne satisfait pas les conditions permettant de fusionner la TVF et l'agrégation, le système revient à un plan d'exécution sans TVF.
Si une requête non fusionnable utilise le temps de traitement comme attribut temporel, la colonne de temps de traitement est matérialisée et utilisée comme attribut temporel des fenêtres créées. Cela entraîne l'impact du watermark de la table source sur les résultats d'agrégation : les fenêtres peuvent se fermer plus tôt que prévu et les données en retard peuvent être rejetées, comme c'est le cas avec les fenêtres en temps événementiel. Évitez les modèles ci-dessous pour prévenir ce problème.
La TVF de fenêtre et l'instruction d'agrégation ne peuvent pas être fusionnées lorsque l'une des conditions suivantes est remplie :
-
Filtrage ou calcul sur les champs temporels de la fenêtre.
window_start,window_endouwindow_timeest filtré ou modifié avant l'agrégation.-- Filtering on window_start > SELECT window_start, window_end, item, SUM(price) AS total_price FROM (SELECT item, price, window_start, window_end FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) WHERE window_start >= TIMESTAMP '2020-04-15 08:06:00.000') GROUP BY item, window_start, window_end; -- Arithmetic on window_start > SELECT window_start, window_end, item, SUM(price) AS total_price FROM (SELECT item, price, window_start + (INTERVAL '1' SECOND) AS window_start, window_end FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))) GROUP BY item, window_start, window_end; -- Type casting on window_start > SELECT window_start, window_end, item, SUM(price) AS total_price FROM (SELECT item, price, CAST(window_start AS varchar) AS window_start, window_end FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))) GROUP BY item, window_start, window_end; -
Une TVF de fenêtre est utilisée avec une fonction table définie par l'utilisateur (UDTF).
> SELECT window_start, window_end, category, SUM(price) AS total_price FROM (SELECT category, price, window_start, window_end FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)), LATERAL TABLE(category_udtf(item)) AS T(category)) GROUP BY category, window_start, window_end; -
La clause
GROUP BYometwindow_startouwindow_end.> SELECT window_start, item, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) GROUP BY item, window_start; Une fonction d'agrégation définie par l'utilisateur Python (UDAF) est utilisée.
-
GROUPING SETS,CUBEouROLLUPeffectuent un regroupement séparé parwindow_startouwindow_end.> SELECT item, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) GROUP BY GROUPING SETS((item), (window_start), (window_end)); > SELECT item, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) GROUP BY CUBE (item, window_start, window_end); > SELECT item, SUM(price) AS total_price FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) GROUP BY ROLLUP (item, window_start, window_end); -
Une fonction d'agrégation est appliquée à
window_start,window_endouwindow_time.> SELECT window_start, window_end, item, SUM(price) AS total_price, MAX(window_end) AS max_end FROM TABLE( SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)) GROUP BY item, window_start, window_end;
Prise en charge des flux de mise à jour
Le tableau ci-dessous indique la prise en charge des flux de mise à jour selon la fonction de fenêtre et la syntaxe.
| Fonction de fenêtre | Ancienne syntaxe (GroupWindowAggregation) — VVR | Ancienne syntaxe (GroupWindowAggregation) — Apache Flink | Nouvelle syntaxe (WindowAggregate) — VVR | Nouvelle syntaxe (WindowAggregate) — Apache Flink |
|---|---|---|---|---|
| TUMBLE | Oui | Oui | Oui | Non |
| HOP | Oui | Oui | Oui | Non |
| SESSION | Oui | Oui | Oui | Oui (Apache Flink 1.19 et ultérieur) |
| CUMULATE | N/A | N/A | Oui (VVR 8.0.6 et ultérieur) | Non |
Avec l'ancienne syntaxe, la prise en charge des flux de mise à jour est identique que vous utilisiez VVR ou Apache Flink. Avec la nouvelle syntaxe, seul l'opérateur WindowAggregate de VVR prend en charge les flux de mise à jour pour toutes les fonctions de fenêtre. VVR sélectionne automatiquement les opérateurs GroupWindowAggregation et WindowAggregate en fonction du flux d'entrée.
Pour connaître les différences entre la fonction de fenêtre SESSION dans VVR et Apache Flink, consultez Requêtes .
Optimisation Local-Global pour l'agrégation par fenêtre
L'agrégation par fenêtre prend en charge l'optimisation d'agrégation en deux phases Local-Global. Lorsqu'elle est activée, l'optimiseur divise une agrégation par fenêtre en une seule phase en :
Local Aggregate : effectue une pré-agrégation partielle avant le brassage des données, réduisant ainsi le volume de données transférées sur le réseau.
Global Aggregate : effectue l'agrégation finale après le brassage et produit les résultats.
Les six conditions suivantes doivent toutes être réunies pour que cette optimisation prenne effet.
Condition 1 : La stratégie de phase d'agrégation autorise la double phase
table.optimizer.agg-phase-strategy est défini sur AUTO (valeur par défaut) ou TWO_PHASE. La définition sur ONE_PHASE désactive l'optimisation en deux phases.
Condition 2 : La fenêtre utilise le temps événementiel
La fenêtre doit utiliser le temps événementiel (rowtime). Les fenêtres en temps de traitement ne sont pas prises en charge.
Condition 3 : Le type de fenêtre n'est pas SESSION
Les types de fenêtre TUMBLE, HOP et CUMULATE sont pris en charge. Les fenêtres SESSION ne prennent pas en charge l'optimisation en deux phases.
Condition 4 : Toutes les fonctions d'agrégation prennent en charge la fusion partielle
Toutes les fonctions d'agrégation doivent prendre en charge l'opération de fusion. Les fonctions intégrées telles que SUM, COUNT, MIN, MAX et AVG sont prises en charge. Les UDAF personnalisées doivent implémenter la méthode merge().
Condition 5 : Le flux d'entrée est en insertion seule et la fenêtre peut être convertie au format TVF
Les conditions suivantes doivent toutes être remplies :
Le flux d'entrée est en insertion seule.
table.exec.emit.early-fire.enabledest défini surfalse(valeur par défaut).table.exec.emit.late-fire.enabledest défini surfalse(valeur par défaut).Pour les fenêtres HOP, la fenêtre doit être alignée (la taille de la fenêtre est divisible par l'intervalle de glissement).
Condition 6 : La distribution des données ne satisfait pas déjà aux exigences de partitionnement
La distribution des données d'entrée ne satisfait pas déjà aux exigences de partitionnement pour l'agrégation. Si les données sont déjà distribuées selon la clé de partitionnement, l'optimiseur détermine qu'aucune pré-agrégation supplémentaire n'est nécessaire et ne génère pas de nœud Local Aggregate.
Référence des paramètres
| Paramètre | Type | Valeur par défaut | Exigence |
|---|---|---|---|
table.optimizer.agg-phase-strategy |
Enum | AUTO | Ne doit pas être ONE_PHASE |
table.exec.emit.early-fire.enabled |
Booléen | false | Doit être false |
table.exec.emit.late-fire.enabled |
Booléen | false | Doit être false |