Cette rubrique répond aux questions courantes concernant la validité des données dans Realtime Compute for Apache Flink.
Pourquoi n'y a-t-il aucune sortie dans la table de destination ?
Comment résoudre l'absence de sortie dans le système en aval ?
Pourquoi n'y a-t-il aucune sortie dans la table de destination ?
Après le démarrage d'un job, aucune donnée n'apparaît dans la table de destination. Effectuez les vérifications suivantes dans l'ordre.

Vérifiez les basculements. Si le job a subi un basculement, analysez le message d'erreur pour identifier la cause racine et la résoudre afin que le job s'exécute comme prévu.
Assurez-vous que les données atteignent Realtime Compute for Apache Flink. En l'absence de basculement mais en cas de latence élevée des données, consultez la métrique
numRecordsInOfSourcesur la page de surveillance et d'alertes. Si cette métrique affiche zéro pour une source, la table source n'envoie pas de données à Flink ; vérifiez la source de données en amont.Vérifiez si un opérateur filtre tous les enregistrements. Ajoutez
pipeline.operator-chaining: 'false'au champ Other Configuration (voir Comment configurer des paramètres d'exécution personnalisés pour un job ?). Cette action divise la chaîne d'opérateurs, ce qui vous permet d'inspecter individuellement les métriques Bytes Received et Bytes Sent de chaque opérateur. Un opérateur recevant des entrées mais produisant zéro sortie est responsable du problème ; les coupables fréquents sont les opérateurs JOIN, WINDOW et WHERE.-
Vérifiez si la base de données en aval conserve les données dans son tampon d'écriture. Réduisez la taille du lot du connecteur en aval pour vider les données plus rapidement.
ImportantÉvitez une taille de lot excessivement petite. Une taille de lot de 1 implique que Flink envoie une requête distincte pour chaque enregistrement traité, ce qui peut surcharger la base de données en aval en cas de volumes de données élevés.
Vérifiez la présence d'interblocages dans ApsaraDB RDS for MySQL. Consultez la section Interblocage lors de l'écriture dans MySQL via le connecteur ApsaraDB RDS ou TDDL.
Pour isoler le problème, imprimez les résultats intermédiaires dans les journaux à l'aide d'une table de destination print. Consultez la section Afficher la sortie du connecteur print .
Absence de sortie due à l'activation de MiniBatch avec consommation partielle du journal binaire en mode CDC
-
Symptôme
Un job CDC commence à consommer le journal binaire à partir du milieu en utilisant
latestou un décalage spécifique, avec MiniBatch activé. La table de destination reçoit alors des données partielles ou aucune sortie. -
Cause
MiniBatch fusionne et annule les messages de journal des modifications pour la même clé primaire au sein d'un lot. Lorsque le job démarre au milieu du journal binaire, il peut recevoir un message
UPDATE_AFTERsans le messageUPDATE_BEFOREcorrespondant (ou inversement). Les messages opposés s'annulent mutuellement dans le lot et ne produisent aucune sortie, ce qui entraîne des données incorrectes en aval. -
Solutions
Consommez le journal binaire depuis le début. Évitez les modes de consommation partielle tels que
latestou un décalage spécifié afin de préserver la séquence complète des modifications pour chaque enregistrement.Désactivez MiniBatch pour que le job traite les enregistrements un par un. Cette option sacrifie une partie du débit au profit de l'exactitude.
Comment résoudre les problèmes de lecture de source Flink ?
Si Realtime Compute for Apache Flink ne parvient pas à lire depuis une source, effectuez les vérifications suivantes.
Connectivité réseau
Par défaut, Realtime Compute for Apache Flink ne peut atteindre que les services situés dans la même région et le même cloud privé virtuel (VPC). Pour un accès inter-réseaux :
Inter-VPC : Comment accéder à d'autres services via des VPC ?
Accès Internet : Comment accéder à Internet ?
Liste d'autorisation du service en amont
Pour lire depuis des services tels que Kafka et Elasticsearch, ajoutez votre espace de travail Flink à leurs listes d'autorisation :
Obtenez le bloc CIDR du vSwitch de votre espace de travail Flink. Consultez la section Comment configurer une liste d'autorisation ?
Ajoutez ce bloc CIDR à la liste d'autorisation du service en amont. Reportez-vous à la section « Prérequis » de la documentation du connecteur, par exemple Kafka.
Cohérence des champs entre la table Flink et la table physique
Les incompatibilités dans les définitions des champs sont une cause fréquente d'échecs de lecture. Lors de la rédaction du DDL pour votre table source Flink :
Ordre des champs : Respectez l'ordre exact des champs de la table physique.
Casse des noms de champs : Utilisez la même casse que dans la table physique.
Type de champ : Utilisez les types équivalents mappés. Vérifiez la section « Mappages des types de données » de la documentation du connecteur concerné, par exemple Simple Log Service.
Exceptions dans les journaux TaskManager
Vérifiez si le journal TaskManager de la table source contient des messages d'exception :
Dans le volet de navigation de gauche, accédez à O&M > Deployment.
Cliquez sur le nom du déploiement.
Cliquez sur l'onglet Status, puis cliquez sur le sommet source dans le DAG.
Dans le panneau de droite, cliquez sur l'onglet SubTasks.
Dans la colonne More, cliquez sur l'icône
et choisissez Open TaskManager Log Page.Sur l'onglet Logs, recherchez la première entrée contenant « Caused by » ; elle indique généralement la cause racine.
Comment résoudre l'absence de sortie dans le système en aval ?
Effectuez les vérifications suivantes.
Connectivité réseau
Par défaut, Realtime Compute for Apache Flink ne peut atteindre que les services situés dans la même région et le même VPC. Pour un accès inter-réseaux :
Inter-VPC : Comment accéder à d'autres services via des VPC ?
Accès Internet : Comment accéder à Internet ?
Liste d'autorisation du système en aval
Pour écrire dans des services tels que ApsaraDB RDS for MySQL, Kafka, Elasticsearch, AnalyticDB for MySQL 3.0, Apache HBase, Redis et ClickHouse, ajoutez votre espace de travail Flink à leurs listes d'autorisation :
Obtenez le bloc CIDR du vSwitch de votre espace de travail Flink. Consultez la section Comment configurer une liste d'autorisation ?
Ajoutez ce bloc CIDR à la liste d'autorisation du service en aval. Reportez-vous à la section « Prérequis » de la documentation du connecteur, par exemple ApsaraDB RDS for MySQL.
Cohérence des champs entre la table Flink et la table physique
Appliquez les mêmes vérifications que celles décrites dans la section Comment résoudre les problèmes de lecture de source Flink ? : vérifiez l'ordre des champs, la casse des noms de champs et les mappages des types de champs.
Données filtrées par les opérateurs
Examinez les comptes d'entrée et de sortie de chaque sommet dans le DAG du job. Si un sommet tel que WHERE affiche entrée = 5 et sortie = 0, cet opérateur rejette tous les enregistrements.
Seuils de tampon du connecteur de destination trop élevés
Lorsque le volume d'entrée est faible, des seuils de tampon par défaut élevés peuvent empêcher le vidage des données vers le système en aval, car le tampon ne se remplit jamais suffisamment pour déclencher une écriture. Réduisez les options pertinentes selon vos besoins :
| Option | Description | Service en aval concerné |
|---|---|---|
batchSize |
Taille des données écrites à la fois | DataHub, Tablestore, MongoDB, ApsaraDB RDS for MySQL, AnalyticDB for MySQL V3.0, ApsaraDB for ClickHouse, TSDB for InfluxDB |
batchCount |
Nombre maximal d'enregistrements écrits à la fois | DataHub |
flushIntervalMs |
Intervalle de vidage du tampon MaxCompute Tunnel Writer | MaxCompute |
sink.buffer-flush.max-size |
Taille des données mises en mémoire tampon avant l'écriture dans HBase, en octets | ApsaraDB for Hbase |
sink.buffer-flush.max-rows |
Nombre d'enregistrements mis en mémoire tampon avant l'écriture dans HBase | ApsaraDB for Hbase |
sink.buffer-flush.interval |
Intervalle auquel les données mises en tampon sont vidées périodiquement vers HBase | ApsaraDB for Hbase |
jdbcWriteBatchSize |
Nombre maximal de lignes traitées par un nœud de destination de flux Hologres à la fois lors de l'utilisation d'un pilote JDBC | Hologres |
Données désordonnées dans les fenêtres basées sur le temps événementiel
Les filigranes contrôlent les enregistrements acceptés par une fenêtre. Si le premier enregistrement possède un horodatage de 2100 et définit le filigrane à 2100, tout enregistrement ultérieur dont l'horodatage est inférieur à 2100 (par exemple 2021) est considéré comme tardif et rejeté. La fenêtre ne peut pas se fermer tant qu'un enregistrement avec un horodatage supérieur à 2100 n'arrive pas.
Pour détecter les enregistrements désordonnés, utilisez une table de destination print ou examinez les journaux Log4j. Consultez les sections Créer une table de destination print et Configurer la sortie des journaux. Si des enregistrements tardifs sont confirmés, filtrez-les ou configurez votre stratégie de filigrane pour autoriser une période de grâce pour les arrivées tardives.
Sous-tâches source sans entrée
Lorsqu'une sous-tâche source ne reçoit aucune donnée, son filigrane reste à la valeur par défaut de l'époque (1970-01-01T00:00:00Z), qui devient le filigrane global de l'opérateur. Cela empêche les fenêtres basées sur le temps événementiel de se fermer.
Vérifiez le DAG du job et confirmez que toutes les sous-tâches source reçoivent des entrées. Si une sous-tâche est inactive, réduisez le parallélisme du job pour qu'il corresponde au nombre de shards de la table en amont, afin que chaque sous-tâche reçoive des données.
Partitions Kafka vides
Une partition Kafka vide peut bloquer la génération de filigranes. Consultez la section Pourquoi une fenêtre de temps événementiel ne produit-elle aucune sortie à partir d'une table source Kafka ?
Comment résoudre les pertes de données ?
Les réductions de volume de données proviennent généralement des clauses WHERE, des jointures ou des opérations fenêtrées. Pour les pertes inexpliquées, effectuez les vérifications suivantes.
Politique de cache de la table de dimension
Une politique de cache incorrecte peut entraîner des échecs de jointure par recherche qui suppriment silencieusement des enregistrements. Configurez une politique de cache appropriée à l'aide des options liées au cache dans la documentation du connecteur, par exemple la section « Spécifique aux tables de dimension (telles que les paramètres Cache) » dans ApsaraDB for Hbase.
Utilisation des fonctions
Une utilisation incorrecte de fonctions telles que to_timestamp_tz et date_format peut provoquer des échecs de conversion de données qui rejettent silencieusement des enregistrements. Vérifiez le comportement des fonctions à l'aide d'une table de destination print ou des journaux Log4j. Consultez les sections Print et Configurer la sortie des journaux.
Données désordonnées
Les événements tardifs sont rejetés lorsque leur horodatage se situe en dehors de la plage acceptée de la fenêtre actuelle. Par exemple, un événement avec un horodatage de 11 s entrant dans une fenêtre de 15–20 s est rejeté car son filigrane est de 11, soit inférieur à la limite inférieure de la fenêtre.

Les pertes dues à cette cause se concentrent généralement dans une seule fenêtre. Utilisez une table de destination print ou Log4j pour confirmer la présence de données désordonnées.
Pour minimiser les pertes dues au désordre, définissez une stratégie de génération de filigrane avec une période de grâce (par exemple, Watermark = Event time - 5s). Alignez les fenêtres sur des limites exactes de jour, d'heure ou de minute ; cela rend le comportement des fenêtres prévisible et réduit les rejets de cas limites lorsqu'il est combiné avec une période de grâce appropriée.
Pourquoi obtenez-vous des résultats inexacts lors de la déduplication des données ingérées depuis Hologres en mode CDC avec ROW_NUMBER ?
SELECT
hg_binlog_timestamp_us,
order_id,
order_name,
order_pay,
order_starttime
FROM(
SELECT
hg_binlog_lsn,
hg_binlog_event_type,
hg_binlog_timestamp_us,
order_id,
order_name,
order_pay,
order_starttime,
ROW_NUMBER()
OVER(
PARTITION BY
order_id
ORDER BY
order_starttime desc
) as rn
FROM test_cdc
)
where rn = 1
Supposons que la table test_cdc contient les exemples de données suivants, dans lesquels la commande avec order_id 1001 possède deux enregistrements de modification :
hg_binlog_timestamp_us | order_id | order_name | order_pay | order_starttime
1785578400000000 | 1001 | phone | 3999.00 | 2026-08-01 10:00:00
1785582000000000 | 1001 | phone | 3599.00 | 2026-08-01 11:00:00
1785578460000000 | 1002 | earphone | 299.00 | 2026-08-01 10:01:00
Après avoir exécuté l'instruction SQL précédente, les enregistrements devraient être partitionnés par order_id et triés par order_starttime dans l'ordre décroissant, et seuls les enregistrements avec rn = 1 devraient être conservés. Cependant, dans ce scénario anormal, la déduplication ne prend pas effet et le résultat contient toujours les deux enregistrements de modification de la commande avec order_id 1001 :
hg_binlog_timestamp_us | order_id | order_name | order_pay | order_starttime
1785578400000000 | 1001 | phone | 3999.00 | 2026-08-01 10:00:00
1785582000000000 | 1001 | phone | 3599.00 | 2026-08-01 11:00:00
1785578460000000 | 1002 | earphone | 299.00 | 2026-08-01 10:01:00
Le système en aval utilise un opérateur de rétraction (par exemple, ROW_NUMBER OVER WINDOW pour la déduplication), mais la source Hologres n'est pas configurée pour émettre des données en mode upsert. Sans le mode upsert, la source émet des événements d'insertion uniquement que l'opérateur de rétraction ne peut pas traiter correctement.
Ajoutez 'upsertSource' = 'true' à la clause WITH de l'instruction DDL de la table source.
,'binlog' = 'true'
,'sdkMode' = 'jdbc'
,'cdcMode' = 'true'
,'binlogStartUpMode' = 'initial'
,'jdbcBinlogSlotName' = 'test_cdc_1_replication_slot'
,'binlogMaxRetryTimes' = '10',
'binlogRetryIntervalMs' = '500',
'binlogBatchReadSize' = '100'
,'upsertSource' = 'true'
)
Comment résoudre les résultats inexacts ?
Activez le profilage des opérateurs pour inspecter les résultats intermédiaires sans modifier la logique du job.
-
Analysez les journaux d'exécution :

Cliquez sur le nom du déploiement, puis cliquez sur l'onglet Status.
Dans le DAG, copiez le nom de l'opérateur produisant des résultats incorrects.
Dans la liste des journaux, cliquez sur
inspect-taskmanager_0.outsous Log Name et recherchez le nom de l'opérateur.
Après avoir identifié la cause racine, révisez la logique de l'opérateur, redémarrez le job et vérifiez l'exactitude des données.
Comment corriger l'erreur « doesn't support consuming update and delete changes which is produced by node TableSourceScan » ?
Le message d'erreur ressemble à ceci :
Table sink 'vvp.default.***' doesn't support consuming update and delete changes which is produced by node TableSourceScan(table=[[vvp, default, ***]], fields=[id,b, content])
at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.wrapExecutor(DelegateOperationExecutor.java:286)
at org.apache.flink.table.sqlserver.execution.DelegateOperationExecutor.validate(DelegateOperationExecutor.java:211)
at org.apache.flink.table.sqlserver.FlinkSqlServiceImpl.validate(FlinkSqlServiceImpl.java:741)
at org.apache.flink.table.sqlserver.proto.FlinkSqlServiceGrpc$MethodHandlers.invoke(FlinkSqlServiceGrpc.java:2522)
at io.grpc.stub.ServerCalls$UnaryServerCallHandler$UnaryServerCallListener.onHalfClose(ServerCalls.java:172)
at io.grpc.internal.ServerCallImpl$ServerStreamListenerImpl.halfClosed(ServerCallImpl.java:331)
at io.grpc.internal.ServerImpl$JumpToApplicationThreadServerStreamListener$1HalfClosed.runInContext(ServerImpl.java:820)
at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:37)
at io.grpc.internal.SerializingExecutor.run(SerializingExecutor.java:123)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1147)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:622)
at java.lang.Thread.run(Thread.java:834)
La table de destination est en mode ajout uniquement et ne peut pas consommer les événements de mise à jour ou de suppression provenant de la source. Remplacez-la par une destination qui prend en charge les upserts, telle que Upsert Kafka.
Comment corriger les écrasements ou suppressions inattendus de données lors de l'utilisation du connecteur Lindorm ?
Par défaut, le connecteur Lindorm utilise l'opérateur upsert materialize (par défaut : AUTO) pour gérer l'ordre d'écriture. Cet opérateur génère une opération DELETE suivie d'une opération INSERT pour la même clé primaire. Deux caractéristiques de Lindorm rendent cela problématique :
Précision de l'horodatage en millisecondes : Lindorm versionne les données à l'aide d'horodatages en millisecondes. Plusieurs enregistrements avec la même clé primaire écrits dans la même milliseconde peuvent arriver dans le désordre, provoquant des conflits de version.
Absence de prise en charge native de DELETE : Lindorm prend uniquement en charge la sémantique UPSERT ; les suppressions sont irréversibles. La logique de maintien de l'ordre de
upsert materializeest donc inefficace et peut provoquer des anomalies de données dues à la séquence DELETE + INSERT.
Lorsque des écritures simultanées aboutissent dans la même milliseconde, les opérations DELETE et INSERT résultantes peuvent produire des données incorrectes ou une perte silencieuse de données.
Solution : Désactivez explicitement l'opérateur upsert materialize en ajoutant ce qui suit à la configuration des paramètres d'exécution de votre job ou au code SQL :
SET 'table.exec.sink.upsert-materialize' = 'NONE';
Ce paramètre s'applique à tout job qui écrit dans Lindorm via Flink.
Après la désactivation de cet opérateur, seule la cohérence à terme est garantie. Confirmez que la cohérence à terme est acceptable pour votre cas d'utilisation avant d'appliquer cette modification.