Lorsque vous utilisez Flink SQL pour le traitement de données en temps réel, des événements de journal des modifications (changelog) désordonnés peuvent corrompre silencieusement les résultats : des enregistrements sont supprimés alors qu'ils devraient exister, ou des mises à jour apparaissent dans le mauvais ordre. Cette rubrique explique pourquoi ces événements surviennent, comment l'opérateur SinkUpsertMaterializer les corrige et comment ajuster ou éviter cet opérateur lorsque les performances sont critiques.
Concepts clés
Journal des modifications et types de flux
Dans les bases de données relationnelles telles que MySQL, le journal binaire (binlog) capture chaque opération INSERT, UPDATE et DELETE. Flink SQL utilise un mécanisme similaire appelé journal des modifications (changelog) pour suivre les changements de données et permettre un traitement incrémentiel dans les pipelines de streaming.
Un flux de journal des modifications appartient à l'une des deux catégories suivantes :
| Type de flux | Types d'événements | Description |
|---|---|---|
| Flux append-only | +I uniquement | Contient uniquement des événements INSERT. Aucune mise à jour ni suppression. Également appelé flux sans mise à jour. |
| Flux de mise à jour | +I, +U, -U, -D | Contient des événements de mise à jour ou de suppression en plus des insertions. Des opérateurs tels que l'agrégation par groupe et la déduplication produisent ce type de flux. |
Tous les opérateurs ne peuvent pas consommer des flux de mise à jour. Les opérateurs d'agrégation globale (over aggregation) et de jointure par intervalle n'acceptent que des flux append-only en entrée.
Types d'événements du journal des modifications
Flink SQL utilise quatre types d'événements, basés sur l'énumération RowKind de l'API Apache Flink :
| Nom court | Nom complet | Sémantique |
|---|---|---|
+I |
INSERT | Insère une nouvelle ligne. |
-U |
UPDATE_BEFORE | Rétracte le contenu précédent d'une ligne mise à jour. Toujours associé à un événement +U. |
+U |
UPDATE_AFTER | Contient le nouveau contenu d'une ligne mise à jour. Toujours associé à un événement -U. |
-D |
DELETE | Supprime une ligne. |
Flink conserve UPDATE_BEFORE (-U) et UPDATE_AFTER (+U) comme types d'événements distincts plutôt que de les combiner en un seul événement UPDATE composite, et ce pour deux raisons :
Structure uniforme : Les deux événements partagent la même structure de ligne, différenciée uniquement par la propriété
RowKind. Un type d'événement composite nécessiterait des structures hétérogènes ou un alignement spécial entre les événements INSERT et DELETE.Brassage distribué : Dans les pipelines parallèles, les opérations de jointure et d'agrégation brassent les données entre les tâches. Les événements UPDATE composites devraient tout de même être divisés en événements séparés lors du brassage pour maintenir la cohérence ; les garder séparés dès le départ simplifie donc le modèle.
Comment surviennent les événements désordonnés
Prenons l'exemple suivant, utilisé tout au long de cette rubrique pour illustrer le problème et la solution :
-- CDC source tables
CREATE TEMPORARY TABLE s1 (
id BIGINT,
level BIGINT,
PRIMARY KEY(id) NOT ENFORCED
) WITH (...);
CREATE TEMPORARY TABLE s2 (
id BIGINT,
attr VARCHAR,
PRIMARY KEY(id) NOT ENFORCED
) WITH (...);
-- Sink table
CREATE TEMPORARY TABLE t1 (
id BIGINT,
level BIGINT,
attr VARCHAR,
PRIMARY KEY(id) NOT ENFORCED
) WITH (...);
-- Join s1 and s2 and write the result to t1
INSERT INTO t1
SELECT s1.*, s2.attr
FROM s1 JOIN s2
ON s1.level = s2.id;
Lorsque l'enregistrement (id=1, level=10) de la table s1 est inséré à l'instant t0, puis mis à jour vers (id=1, level=20) à l'instant t1, trois événements de journal des modifications sont produits :
| Événement | Type |
|---|---|
+I (id=1, level=10) |
INSERT |
-U (id=1, level=10) |
UPDATE_BEFORE |
+U (id=1, level=20) |
UPDATE_AFTER |
La clé primaire de s1 est id, mais la clause JOIN brasse les données sur la colonne level. Avec un parallélisme de 2 pour l'opérateur Join, ces trois événements peuvent être acheminés vers deux tâches différentes — l'une gérant level=10 et l'autre level=20.


Comme les événements sont traités en parallèle, l'opérateur Sink en aval peut les recevoir selon l'un des trois ordres possibles suivants :
| Cas 1 (ordre correct) | Cas 2 (désordonné) | Cas 3 (désordonné) |
|---|---|---|
+I (id=1, level=10, attr='a1') |
+U (id=1, level=20, attr='b1') |
+I (id=1, level=10, attr='a1') |
-U (id=1, level=10, attr='a1') |
+I (id=1, level=10, attr='a1') |
+U (id=1, level=20, attr='b1') |
+U (id=1, level=20, attr='b1') |
-U (id=1, level=10, attr='a1') |
-U (id=1, level=10, attr='a1') |

Le Cas 1 traite les événements dans leur séquence d'origine : aucun problème. Dans les Cas 2 et 3, la table sink a id comme clé primaire. Si le stockage externe effectue un upsert, l'enregistrement avec id=1 finit par être supprimé, alors que l'état final attendu est (id=1, level=20, attr='b1').
Les événements désordonnés ne se produisent que lorsque le parallélisme de l'opérateur Join est supérieur à 1. Une paire d'événements partageant la même clé upsert est toujours acheminée vers la même tâche, ce qui explique pourquoi seuls trois cas d'ordonnancement sont possibles dans ce scénario.
SinkUpsertMaterializer
Fonctionnement
SinkUpsertMaterializer est un opérateur intermédiaire que Flink insère pour résoudre les problèmes d'ordonnancement. Il a été introduit pour traiter le ticket FLINK-20374.
Pour comprendre pourquoi SinkUpsertMaterializer est nécessaire, il faut appréhender la notion de clé upsert. Une upsert key est une colonne (ou un ensemble de colonnes) qui préserve l'ordre de tri d'une clé unique à travers une opération SQL. Lorsque des clés upsert existent, l'opérateur en aval reçoit les événements de mise à jour dans le bon ordre. Lorsqu'une opération de brassage de données rompt l'ordre de la clé unique — comme le fait la jointure sur level dans cet exemple — la clé upsert devient vide.
Dans cet exemple, les lignes de s1 sont brassées par level, de sorte que la sortie du Join contient des lignes ayant la même valeur s1.id mais dans un ordre arbitraire. Les clés uniques sont (s1.id), (s1.id, s1.level) et (s1.id, s2.id), mais la clé upsert est vide. De plus, la clé primaire de la table sink (id) ne correspond pas à la clé upsert de la sortie du Join. SinkUpsertMaterializer comble cet écart.
Les événements de journal des modifications désordonnés suivent des règles spécifiques : pour une clé upsert donnée (ou pour toutes les colonnes si la clé upsert est vide), les événements ADD (+I et +U) surviennent toujours avant les événements RETRACT correspondants (-D et -U). Une paire d'événements de journal des modifications partageant la même clé upsert est traitée par la même tâche, même en cas de brassage de données. Ce sont ces garanties d'ordonnancement sur lesquelles s'appuie SinkUpsertMaterializer pour reconstruire des résultats corrects.
L'opérateur fonctionne comme suit :
Il maintient une liste de valeurs
RowDatadans l'état, indexée par la clé upsert déduite (ou par la ligne entière si la clé upsert est vide).Lors d'un événement ADD (
+Iou+U) : il ajoute ou met à jour la ligne dans l'état.Lors d'un événement RETRACT (
-Uou-D) : il supprime la ligne de l'état.Il génère des événements de journal des modifications corrects en fonction de la clé primaire de la table sink.

Le schéma suivant montre comment SinkUpsertMaterializer gère les Cas 2 et 3 de l'exemple ci-dessus :
Cas 2 : Lorsque
-U (id=1, level=10, attr='a1')arrive en dernier, SinkUpsertMaterializer supprime cette ligne de l'état et génère un événement UPDATE basé sur l'avant-dernière ligne. Le résultat final est(id=1, level=20, attr='b1').Cas 3 : Lorsque
+U (id=1, level=20, attr='b1')arrive, l'opérateur le transmet en aval. Lorsque-U (id=1, level=10, attr='a1')arrive plus tard, l'opérateur supprime la ligne correspondante de l'état sans émettre d'événement. Le résultat final est encore(id=1, level=20, attr='b1').

Pour consulter le code source, voir SinkUpsertMaterializer (Flink release-1,17).
Quand SinkUpsertMaterializer est déclenché
Flink ajoute l'opérateur SinkUpsertMaterializer dans les scénarios suivants :
-
La table sink possède une clé primaire, mais les données entrantes ne satisfont pas la contrainte UNIQUE. Les causes courantes incluent :
Définir une clé primaire sur la table sink alors que la table source n'en a pas.
Exclure la colonne de clé primaire source lors de l'écriture dans le sink, ou mapper une colonne non-clé primaire de la source vers la clé primaire du sink.
Réduire la précision d'une colonne de clé primaire via une conversion de type ou une agrégation par groupe (par exemple, un cast de BIGINT vers INT).
-
Transformer la colonne de clé primaire, comme concaténer plusieurs colonnes en une seule :
CREATE TABLE students ( student_id BIGINT NOT NULL, student_name STRING NOT NULL, course_id BIGINT NOT NULL, score DOUBLE NOT NULL, PRIMARY KEY(student_id) NOT ENFORCED ) WITH (...); CREATE TABLE performance_report ( student_info STRING NOT NULL PRIMARY KEY NOT ENFORCED, avg_score DOUBLE NOT NULL ) WITH (...); CREATE TEMPORARY VIEW v AS SELECT student_id, student_name, AVG(score) AS avg_score FROM students GROUP BY student_id, student_name; -- The concatenated result no longer satisfies the UNIQUE constraint -- but is used as the primary key of the sink table. INSERT INTO performance_report SELECT CONCAT('id:', student_id, ',name:', student_name) AS student_info, avg_score FROM v;
Une opération de brassage de données perturbe l'ordre de tri d'une clé unique avant l'écriture dans la table sink. C'est le scénario de l'exemple de jointure ci-dessus : la jointure sur
levelbrasse les lignes de s1, rompant l'ordre de tri de la clé primaireid.Le paramètre
table.exec.sink.upsert-materializeest défini surforce.
Configurer SinkUpsertMaterializer
Utilisez le paramètre table.exec.sink.upsert-materialize pour contrôler quand Flink ajoute l'opérateur SinkUpsertMaterializer :
| Valeur | Comportement |
|---|---|
auto (par défaut) |
Flink déduit si des événements désordonnés sont possibles et ajoute l'opérateur si nécessaire. |
none |
Désactive complètement l'opérateur. |
force |
Ajoute toujours l'opérateur, même lorsqu'aucune clé primaire n'est définie sur la table sink. |
Définirautone garantit pas que les événements sont effectivement désordonnés. Par exemple, utiliser une clauseGROUPING SETSavecCOALESCEpour convertir les valeurs null peut empêcher le planificateur SQL de déterminer si la clé upsert correspond à la clé primaire du sink. Dans ce cas, Flink ajoute SinkUpsertMaterializer par précaution. Si les résultats sont corrects sans l'opérateur, définisseztable.exec.sink.upsert-materializesurnone.
Pour plus d'informations sur les opérations de requête prises en charge dans Realtime Compute for Apache Flink utilisant Ververica Runtime (VVR) 6,0 ou ultérieur, les opérateurs d'exécution correspondants et la prise en charge des flux de mise à jour, consultez Exécution des requêtes.
Notes sur les performances et l'exploitation
SinkUpsertMaterializer maintient un état pour chaque ligne qu'il traite. Cela augmente la taille de l'état et ajoute une surcharge d'E/S pour les lectures et écritures d'état, ce qui réduit le débit. Évitez d'utiliser cet opérateur lorsque c'est possible.
Éviter le déclenchement de SinkUpsertMaterializer
Assurez-vous que la clé de partition utilisée pour la déduplication ou l'agrégation par groupe correspond à la clé primaire de la table sink.
Si un seul niveau de parallélisme convient à votre jeu de données et que vous souhaitez éviter les événements désordonnés, définissez le parallélisme sur 1 et désactivez SinkUpsertMaterializer en définissant
table.exec.sink.upsert-materializesurnone.S'il existe une chaîne d'opérateurs entre l'opérateur Sink et un opérateur avec état en amont (tel qu'un opérateur de déduplication ou d'agrégation par groupe), et qu'aucun problème de précision des données ne s'est produit avec les versions de VVR antérieures à 6,0, migrez le déploiement vers VVR 6,0 ou ultérieur. Définissez
table.exec.sink.upsert-materializesurnoneet laissez les autres configurations inchangées. Pour les étapes de migration, consultez Mettre à niveau la version du moteur des déploiements.
Quand vous devez utiliser SinkUpsertMaterializer
N'écrivez pas dans la table sink les colonnes produites par des fonctions non déterministes (telles que
CURRENT_TIMESTAMPouNOW()). Lorsque la clé upsert n'est pas disponible, SinkUpsertMaterializer compare les lignes entières, et les valeurs non déterministes empêchent la correspondance et la suppression des lignes historiques, provoquant une croissance illimitée de l'état.Si l'état de l'opérateur devient suffisamment volumineux pour affecter les performances, augmentez le parallélisme du déploiement. Consultez Configurer les ressources d'un déploiement.
Problèmes connus
SinkUpsertMaterializer peut provoquer une croissance illimitée de l'état dans les situations suivantes :
-
Aucun TTL, TTL trop long ou TTL trop court : Sans durée de vie (TTL) configurée, l'état s'accumule indéfiniment. Un TTL excessivement court peut également poser problème : si l'intervalle entre un événement DELETE et son événement ADD correspondant dépasse le TTL configuré, Flink conserve la ligne dans l'état en tant que données erronées (voir FLINK-29225) et produit le message de journal suivant :
int index = findremoveFirst(values, row); if (index == -1) { LOG.info(STATE_CLEARED_WARN_MSG); return; }Configurez le TTL en fonction de vos besoins métier. Consultez Configurer un déploiement. Realtime Compute for Apache Flink avec VVR 8,0,7 ou ultérieur prend en charge la configuration du TTL par opérateur pour réduire la consommation de ressources des déploiements à grand état. Consultez Configurer le parallélisme, la stratégie de chaînage et le TTL d'un opérateur.
Colonnes non déterministes sans clé upsert : Si le flux de mise à jour arrivant à SinkUpsertMaterializer n'a pas de clé upsert déductible et inclut des colonnes issues de fonctions non déterministes, les lignes historiques ne peuvent pas être mises en correspondance par valeur et ne sont jamais supprimées, entraînant une croissance continue de l'état.
Étapes suivantes
Notes de version — Correspondance des versions du moteur entre Realtime Compute for Apache Flink et Apache Flink
Exécution des requêtes — Opérations de requête prises en charge et prise en charge des flux de mise à jour pour VVR 6,0 et ultérieur
Configurer les ressources d'un déploiement — Augmenter le parallélisme pour gérer un état SinkUpsertMaterializer volumineux