Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Checkpoints et savepoints

Dernière mise à jour :Aug 09, 2026

Cette rubrique répond aux questions fréquentes (FAQ) concernant les checkpoints et les savepoints dans Realtime Compute for Apache Flink.

Échecs de mise à jour des données en mode mini-batch

L'état conserve les résultats des calculs complets précédents. Si le TTL (Time To Live) de l'état expire, celui-ci est effacé et les résultats accumulés sont perdus. Par conséquent, les nouvelles données ne peuvent pas être mises à jour sur la base des résultats du mini-batch.

À l'inverse, si le mini-batch est désactivé, les données associées à la clé expirée sont recalculées et émises lors de l'expiration du TTL de l'état. Cette approche garantit une mise à jour continue des données. Toutefois, l'augmentation de la fréquence des mises à jour peut entraîner d'autres problèmes, tels que des retards de traitement.

Configurez donc les paramètres mini-batch et TTL en fonction de vos besoins métier spécifiques.

Calcul de l'heure de début du prochain checkpoint

L'heure de début du prochain checkpoint dépend de deux paramètres : l'intervalle de checkpoint et la pause minimale entre les checkpoints. Un nouveau checkpoint est déclenché lorsque les deux conditions suivantes sont remplies :

  • Intervalle de checkpoint : durée minimale entre le début d'un checkpoint et le début du suivant. Cela correspond au temps écoulé entre <heure de début du checkpoint précédent, heure de début du checkpoint suivant>.

  • Pause minimale : durée minimale entre la fin d'un checkpoint et le début du suivant. Cela correspond au temps écoulé entre <heure de fin du checkpoint précédent, heure de début du checkpoint suivant>.

Prenons deux scénarios où l'intervalle de checkpoint est de 3 minutes, la pause minimale de 3 minutes et le délai d'expiration de 10 minutes.

  • Scénario 1 : Le déploiement s'exécute normalement et chaque checkpoint réussit.

    Le premier checkpoint commence à 12:00:00 et se termine avec succès à 12:00:02. Le deuxième checkpoint démarre à 12:03:00.

  • Scénario 2 : Le déploiement rencontre des anomalies et un checkpoint échoue en raison d'un délai d'expiration.

    Le premier checkpoint commence à 12:00:00 et se termine avec succès à 12:00:02. Le deuxième checkpoint démarre à 12:03:00 mais échoue à 12:13:00 en raison d'un délai d'expiration. Le troisième checkpoint démarre à 12:16:00.

Pour plus d'informations sur la configuration de la pause minimale entre les checkpoints, consultez Tuning Checkpointing.

GeminiStateBackend dans VVR 8.x par rapport à VVR 6.x

Par défaut, les moteurs Realtime Compute for Apache Flink pour VVR 6.x utilisent GeminiStateBackend V3, tandis que les moteurs VVR 8.x utilisent la version V4.

Catégorie

Description

Fonctionnalités de base

  • V3 (hérité) : prend en charge des fonctionnalités telles que la séparation KV, la séparation calcul-stockage, les savepoints au format standard ou natif, ainsi que le chargement différé de l'état.

  • V4 (nouveau) : l'architecture centrale a été repensée pour le traitement de flux. En plus de prendre en charge toutes les fonctionnalités de V3, V4 offre un accès à l'état et une mise à l'échelle plus rapides.

Paramètre de chargement différé de l'état

  • V4 : state.backend.gemini.file.cache.download.type: LazyDownloadOnRestore

  • V3 : state.backend.gemini.file.cache.lazy-restore: ON

Utilisation de la mémoire gérée

La seule différence réside dans la métrique Resident Set Size (RSS) :

  • V4 : la mémoire n'est allouée par le système d'exploitation qu'en cas de besoin réel, ce qui se reflète alors dans la métrique RSS.

  • V3 : demande directement state's managed memory * 80% au système d'exploitation et gère cette mémoire en interne. Cette allocation apparaît dans la métrique RSS dès le démarrage du déploiement.

Remarque

Pour plus d'informations sur la mémoire gérée, consultez TaskManager Memory.

Checkpoints complets et incrémentiels de même taille

Si vous constatez qu'un checkpoint complet et un checkpoint incrémentiel ont la même taille, procédez comme suit :

  • Vérifiez que la fonction de checkpoint incrémentiel est correctement configurée et activée.

  • Ce comportement peut être attendu dans certains scénarios. Par exemple :

    1. Avant l'ingestion des données (par exemple, avant 18 h 29), le déploiement n'a traité aucune donnée. Le checkpoint ne contient que l'état initial de la source, ce qui en fait effectivement un checkpoint complet.

    2. À 18 h 29, un million d'enregistrements sont ingérés. Si ces données sont entièrement traitées dans l'intervalle du prochain checkpoint (par exemple, 3 minutes) et qu'aucune autre donnée n'arrive, le premier checkpoint incrémentiel contiendra tout l'état généré par ces enregistrements.

    Dans ce cas, il est normal que le checkpoint complet et le premier checkpoint incrémentiel aient la même taille. Le premier checkpoint incrémentiel doit inclure l'état de toutes les données afin de garantir une récupération complète à partir de ce point, ce qui le rend fonctionnellement équivalent à un checkpoint complet.

    Les avantages du checkpoint incrémentiel deviennent généralement apparents à partir du deuxième checkpoint. Avec un flux de données stable et sans modifications majeures de l'état, les checkpoints incrémentiels suivants devraient être plus petits, indiquant que le système capture correctement uniquement les changements d'état. S'ils conservent la même taille, examinez l'état et le comportement de votre système pour identifier d'éventuels problèmes.

Lenteur des checkpoints dans les déploiements Python

  • Cause

    Une fonction définie par l'utilisateur (UDF) Python aux performances médiocres peut augmenter la durée des checkpoints et dégrader les performances du déploiement.

  • Solution

    Réduisez la taille du tampon. Dans la section Other Configuration, définissez les paramètres suivants. Pour les instructions, consultez Configure custom deployment parameters.

    python.fn-execution.bundle.size: Default value: 100000. Unit: records.
    python.fn-execution.bundle.time: Default value: 1000. Unit: milliseconds.

    Pour plus d'informations sur ces paramètres, consultez Flink Python Configuration.

Dépannage des exceptions de checkpoint

  1. Diagnostiquer le type d'exception

    Consultez l'historique des checkpoints dans l'onglet Alarm ou State pour identifier le type d'exception, tel qu'un délai d'expiration ou un échec d'écriture.

    Sélectionnez l'onglet Overview et développez la section Checkpoint. Le tableau affiche les détails de chaque checkpoint, y compris son ID, son statut, son heure de déclenchement, sa durée et la taille des données. La colonne Status indique si chaque checkpoint a réussi ou échoué.

  2. Isoler et résoudre le problème

    • Scénario 1 : Délais d'expiration fréquents des checkpoints. Vérifiez la présence d'une contre-pression (backpressure) dans le déploiement. Analysez la cause racine de la contre-pression, identifiez l'opérateur lent et résolvez le problème en ajustant les ressources ou les configurations. Pour plus d'informations, consultez How to troubleshoot backpressure issues.

    • Scénario 2 : Échecs d'écriture des checkpoints. Suivez ces étapes pour trouver les journaux TaskManager pertinents et analyser la cause racine.

      1. Sur la page Checkpoints de l'onglet Logs, cliquez sur Checkpoints History.

        Sur la page Checkpoints History, vous pouvez consulter les détails de chaque checkpoint, tels que ID, Status, Acknowledged, Trigger Time, End to End Duration et Checkpointed Data Size.

      2. Cliquez sur le signe plus (+) à côté du checkpoint ayant échoué pour afficher les détails de ses opérateurs.

      3. Développez l'opérateur ayant échoué et cliquez sur le SubTask ID pour accéder aux journaux TaskManager correspondants.

Erreur : Restauration d'un ancien état avec le moteur V4

  • Message d'erreur

    Lors de la mise à niveau de VVR 6.x vers VVR 8.x, vous pouvez rencontrer l'erreur suivante : You are using the new V4 state engine to restore old state data from a checkpoint

  • Cause

    VVR 6.x et VVR 8.x utilisent des versions différentes de GeminiStateBackend, et leurs checkpoints ne sont pas compatibles.

  • Solution

    Utilisez l'une des méthodes suivantes pour résoudre ce problème :

    • Créez un savepoint au format standard et démarrez le déploiement à partir de cet état. Pour plus d'informations, consultez Manually create a savepoint et Start a deployment.

    • Redémarrez le déploiement sans état.

    • (Non recommandé) Continuez à utiliser la version héritée de Gemini. Vous devez définir le paramètre state.backend.gemini.engine.type: STREAMING et redémarrer le déploiement pour que la modification prenne effet. Pour savoir comment configurer les paramètres, consultez How to configure deployment parameters.

    • (Non recommandé) Continuez à utiliser le moteur VVR 6.x pour démarrer le déploiement.

Erreur : java.lang.NegativeArraySizeException

  • Message d'erreur

    Un déploiement utilisant un état de liste peut rencontrer l'exception suivante lors de l'exécution :

    Caused by: java.lang.NegativeArraySizeException
      at com.alibaba.gemini.engine.rm.GUnPooledByteBuffer.newTempBuffer(GUnPooledByteBuffer.java:270)
      at com.alibaba.gemini.engine.page.bmap.BinaryValue.merge(BinaryValue.java:85)
      at com.alibaba.gemini.engine.page.bmap.BinaryValue.merge(BinaryValue.java:75)
      at com.alibaba.gemini.engine.pagestore.PageStoreImpl.internalGet(PageStoreImpl.java:428)
      at com.alibaba.gemini.engine.pagestore.PageStoreImpl.get(PageStoreImpl.java:271)
      at com.alibaba.gemini.engine.pagestore.PageStoreImpl.get(PageStoreImpl.java:112)
      at com.alibaba.gemini.engine.table.BinaryKListTable.get(BinaryKListTable.java:118)
      at com.alibaba.gemini.engine.table.BinaryKListTable.get(BinaryKListTable.java:57)
      at com.alibaba.flink.statebackend.gemini.subkeyed.GeminiSubKeyedListStateImpl.getOrDefault(GeminiSubKeyedListStateImpl.java:97)
      at com.alibaba.flink.statebackend.gemini.subkeyed.GeminiSubKeyedListStateImpl.get(GeminiSubKeyedListStateImpl.java:88)
      at com.alibaba.flink.statebackend.gemini.subkeyed.GeminiSubKeyedListStateImpl.get(GeminiSubKeyedListStateImpl.java:47)
      at com.alibaba.flink.statebackend.gemini.context.ContextSubKeyedListState.get(ContextSubKeyedListState.java:60)
      at com.alibaba.flink.statebackend.gemini.context.ContextSubKeyedListState.get(ContextSubKeyedListState.java:44)
      at org.apache.flink.streaming.runtime.operators.windowing.WindowOperator.onProcessingTime(WindowOperator.java:533)
      at org.apache.flink.streaming.api.operators.InternalTimerServiceImpl.onProcessingTime(InternalTimerServiceImpl.java:289)
      at org.apache.flink.streaming.runtime.tasks.StreamTask.invokeProcessingTimeCallback(StreamTask.java:1435)
  • Cause

    Les données d'état pour une seule clé dans l'état de liste ont dépassé 2 Go. Cela peut se produire de la manière suivante :

    1. Pendant le fonctionnement normal, les valeurs ajoutées à une seule clé dans un état de liste sont combinées via un processus de fusion (par exemple, dans un opérateur de fenêtre), ce qui entraîne une croissance continue des données d'état.

    2. Lorsque les données d'état atteignent une certaine taille, elles peuvent d'abord déclencher une erreur d'épuisement de la mémoire (OOM). Après la récupération du déploiement suite à l'échec, le processus de fusion peut amener le backend d'état à demander un tableau d'octets temporaire dépassant 2 Go, provoquant ainsi cette exception.

    Remarque

    RocksDBStateBackend peut rencontrer un problème similaire, qui peut déclencher une ArrayIndexOutOfBoundsException ou une erreur de segmentation. Pour plus d'informations, consultez The EmbeddedRocksDBStateBackend.

  • Solution

    • Si l'état volumineux est causé par un opérateur de fenêtre, envisagez de réduire la taille de la fenêtre.

    • Si l'état volumineux est dû à la logique du déploiement, envisagez de la repenser, par exemple en divisant les clés.

Erreur : FlinkKafkaException: Too many ongoing snapshots

  • Message d'erreur

    org.apache.flink.streaming.connectors.kafka.FlinkKafkaException: Too many ongoing snapshots. Increase kafka producers pool size or decrease number of concurrent checkpoints
  • Cause

    Cette erreur survient lors de l'utilisation d'un sink Kafka et est causée par plusieurs échecs de checkpoint consécutifs.

  • Solution

    Pour éviter les échecs dus aux délais d'expiration, augmentez le délai d'expiration des checkpoints en ajustant le paramètre execution.checkpointing.timeout. Pour plus d'informations sur la configuration des paramètres, consultez Configure custom deployment parameters.

Erreur : Exceeded checkpoint tolerable failure threshold

  • Message d'erreur

    org.apache.flink.util.FlinkRuntimeException:Exceeded checkpoint tolerable failure threshold.
      at org.apache.flink.runtime.checkpoint.CheckpointFailureManager.handleJobLevelCheckpointException(CheckpointFailureManager.java:66)
  • Cause

    Le nombre configuré d'échecs de checkpoint tolérables est trop faible, ce qui provoque un basculement du déploiement lorsque ce seuil est dépassé. Si ce paramètre n'est pas défini, la valeur par défaut est 0, ce qui signifie qu'aucun échec de checkpoint n'est toléré.

  • Solution

    Ajustez le nombre d'échecs de checkpoint autorisés en définissant le paramètre execution.checkpointing.tolerable-failed-checkpoints: num, où num doit être 0 ou un entier positif. Pour plus d'informations sur la configuration des paramètres, consultez Configure custom deployment parameters.