Aller au contenu principal

Points de contrôle d'état asynchrones pour les queries avec état

remarque

Disponible dans Databricks Runtime 10.4 LTS et versions ultérieures.

Le checkpointing d'état asynchrone maintient des garanties d'exécution exactement une fois pour les requêtes de streaming, mais peut réduire la latence globale pour certaines charges de travail stateful de Structured Streaming qui sont limitées par les mises à jour d'état. Ceci est accompli en commençant à traiter le micro-batch suivant dès que le calcul du micro-batch précédent est terminé, sans attendre la fin du checkpointing d'état. Le tableau suivant compare les compromis entre le checkpointing synchrone et asynchrone :

Caractéristique

Point de contrôle synchrone

Point de contrôle asynchrones

Latence

Latence plus élevée pour chaque micro-batch.

Latence réduite, car les micro-batchs peuvent se chevaucher.

Redémarrer

Récupération rapide, car seul le dernier batch doit être réexécuté.

Délai de redémarrage plus élevé car plus d'un micro-batch pourrait devoir être réexécuté.

Caractéristique

Point de contrôle synchrone

Point de contrôle asynchrones

Latence

Latence plus élevée pour chaque micro-batch.

Latence réduite, car les micro-batchs peuvent se chevaucher.

Redémarrer

Récupération rapide, car seul le dernier batch doit être réexécuté.

Délai de redémarrage plus élevé car plus d'un micro-batch pourrait devoir être réexécuté.

Voici les caractéristiques des jobs de streaming qui pourraient bénéficier du point de contrôle d'état asynchrone :

  • Le Job a une ou plusieurs opérations avec état (par exemple, agrégation, flatMapGroupsWithState, mapGroupsWithState, jointures stream-stream).
  • La latence du point de contrôle d’état est l’un des principaux contributeurs à la latence globale d’exécution du batch. Ces informations se trouvent dans les événements StreamingQueryProgress. Ces événements se trouvent également dans les logs log4j du Driver Spark. Voici un exemple de progression de query streaming et comment trouver l’impact du point de contrôle d’état sur la latence globale d’exécution du batch.
    • JSON
       {
      "id" : "2e3495a2-de2c-4a6a-9a8e-f6d4c4796f19",
      "runId" : "e36e9d7e-d2b1-4a43-b0b3-e875e767e1fe",
      "...",
      "batchId" : 0,
      "durationMs" : {
      "...",
      "triggerExecution" : 547730,
      "..."
      },
      "stateOperators" : [ {
      "...",
      "commitTimeMs" : 3186626,
      "numShufflePartitions" : 64,
      "..."
      }]
      }
    • Analyse de la latence du point de contrôle d'état de l'événement de progression de la query ci-dessus

      • La durée du batch (durationMs.triggerDuration) est d’environ 547 secondes.
      • La latence de commit du magasin d'état (stateOperations[0].commitTimeMs) est d'environ 3 186 s. La latence de commit est agrégée sur l'ensemble des tâches contenant un magasin d'état. Dans ce cas, il y a 64 tâches de ce type (stateOperators[0].numShufflePartitions).
      • Chaque tâche contenant un opérateur d'état a pris en moyenne 50 secondes (3 186/64) pour le point de contrôle. Il s'agit d'une latence supplémentaire qui contribue à la durée du batch. En supposant que les 64 tâches s'exécutent simultanément, l'étape de point de contrôle a contribué à environ 9 % (50 secondes / 547 secondes) de la durée du batch. Le pourcentage est encore plus élevé lorsque le nombre maximal de tâches concurrentes est inférieur à 64.

Activer le pointage d'état asynchrone

Vous devez utiliser le magasin d'état basé sur RocksDB pour les points de contrôle d'état asynchrones. Définissez les configurations suivantes :

Scala

spark.conf.set(
"spark.databricks.streaming.statefulOperator.asyncCheckpoint.enabled",
"true"
)

spark.conf.set(
"spark.sql.streaming.stateStore.providerClass",
"com.databricks.sql.streaming.state.RocksDBStateStoreProvider"
)

Limitations et exigences pour le checkpointing asynchrone

remarque

La mise à l’échelle automatique de compute présente des limites pour réduire la taille du cluster pour les charges de travail Structured Streaming. Databricks recommande d'utiliser Spark Declarative Pipelines sur Lakeflow avec une mise à l'échelle automatique améliorée pour les charges de travail de streaming. Voir Optimiser l’utilisation des clusters de LakeFlow Pipelines avec l’autoscaling.

  • Tout échec d'un point de contrôle asynchrone dans un ou plusieurs magasins fait échouer la query. En mode de point de contrôle synchrone, le point de contrôle est exécuté dans le cadre de la tâche, et Spark relance la tâche plusieurs fois avant que la query n'échoue. Ce mécanisme n'est pas présent avec la journalisation asynchrone des états. Databricks recommande l'utilisation de Jobs continus pour les nouvelles tentatives automatiques en cas d'échec de Job. Consultez Exécuter des Jobs en continu.
  • Le checkpointing asynchrone fonctionne mieux lorsque les emplacements du magasin d'état ne sont pas modifiés entre les exécutions de micro-batch. Le redimensionnement de cluster, en combinaison avec le checkpointing d'état asynchrone, pourrait ne pas bien fonctionner car l'instance des magasins d'état pourrait être redistribuée au fur et à mesure que des nœuds sont ajoutés ou supprimés dans le cadre de l'événement de redimensionnement de cluster.
  • Le point de contrôle d'état asynchrone n'est pris en charge que dans l'implémentation du fournisseur de magasin d'état RocksDB. L'implémentation par default du magasin d'états en mémoire ne le prend pas en charge.