Aller au contenu principal

Points de contrôle de Structured Streaming

Les checkpoints et les logs de pré-écriture fonctionnent ensemble pour fournir des garanties de traitement pour les workloads de Structured Streaming. Le checkpoint suit les informations qui identifient la query, y compris les informations d'état et les enregistrements traités. Lorsque vous supprimez les fichiers dans un répertoire de checkpoint ou que vous passez à un nouvel emplacement de checkpoint, la prochaine exécution de la query repart à zéro.

Un répertoire de points de contrôle contient les éléments suivants :

  • **Offsets** : Les offsets sources traités dans chaque micro-batch. Cela permet à la query de reprendre exactement là où elle s'est arrêtée sans retraiter les données.
  • Commits : un enregistrement des micro-batchs qui ont été validés dans le récepteur, ce qui permet la sémantique exactement une fois.
  • État : Pour les query avec état (agrégations, jointures Stream-Stream, déduplication et opérateurs avec état personnalisés comme transformWithState), le point de contrôle stocke les métadonnées sur l'opérateur avec état, le schéma d'état et le contenu du magasin d'état pointé géré par le fournisseur de magasin d'état.
  • Métadonnées : l'ID de query unique utilisé pour identifier la query. Les paramètres de configuration sont stockés dans le log de décalage.

Chaque query doit avoir un emplacement de point de contrôle différent. Plusieurs queries ne devraient jamais partager le même emplacement.

remarque

Cet article couvre les points de contrôle de Structured Streaming pour les query de streaming. Pour des informations sur l’utilisation de DataFrame.checkpoint() avec les volumes Unity Catalog pour tronquer les plans d’exécution des DataFrames non-streaming, consultez points de contrôle des DataFrame dans les volumes.

Activer la création de points de contrôle pour les requêtes Structured Streaming

Vous devez spécifier l'option checkpointLocation avant d'exécuter une query de streaming, comme dans l'exemple suivant :

Python
(df.writeStream
.option("checkpointLocation", "/Volumes/catalog/schema/volume/path")
.toTable("catalog.schema.table")
)
remarque

Certains récepteurs, tels que la sortie pour display() dans les Notebooks et le récepteur memory, génèrent automatiquement un emplacement temporaire de point de contrôle si vous omettez cette option. Ces emplacements de point de contrôle temporaires n'assurent aucune tolérance aux pannes ni garantie de cohérence des données et pourraient ne pas être nettoyés correctement. Databricks recommande de toujours spécifier un emplacement de point de contrôle pour ces récepteurs.

Récupérer après des modifications dans une requête de Structured Streaming

Il existe des limitations concernant les modifications autorisées dans une requête de streaming entre les redémarrages à partir du même emplacement de point de contrôle.

Les modifications qui nécessitent généralement un nouveau point de contrôle incluent le nombre ou le type de sources d'entrée, les rubriques Kafka auxquelles vous êtes abonné ou les chemins Auto Loader, les types d'opérations avec état, le schéma d'état et le type de destination de sortie.

Les modifications généralement sûres comprennent l’ajout ou la suppression de filtres, la modification des limites de débit, des intervalles de Trigger et la mise à jour de la logique de fonction définie par l’utilisateur dans mapGroupsWithState (bien que la sémantique puisse changer).

La section suivante décrit les changements qui ne sont pas autorisés ou dont l'effet n'est pas bien défini, où :

  • Le terme autorisé signifie que vous pouvez effectuer le changement spécifié, mais que la sémantique de son effet soit bien définie dépend de la query et du changement.
  • Le terme *non autorisé* signifie que vous ne devriez pas effectuer la modification spécifiée, car la query redémarrée est susceptible d'échouer avec des erreurs imprévisibles.
  • sdf représente un DataFrame/dataset de streaming généré avec sparkSession.readStream.

Types de modifications dans les queries de Structured Streaming

  • Changements dans le nombre ou le type de sources d'entrée : Ceci n'est pas autorisé par default car Structured Streaming identifie les sources par leur position dans le plan de requête. Si vous activez le nommage des sources, vous pouvez réorganiser les sources existantes et ajouter de nouvelles sources sans repartir d'un nouveau point de contrôle. Consultez Modifier les sources de streaming avec l'évolution des sources.

  • **Modifications des parameters des sources d'entrée** : La possibilité de cela et la bonne définition de la sémantique de la modification dépendent de la source et de la query, y compris les contrôles d'admission tels maxFilesPerTrigger que maxOffsetsPerTrigger ou. Voici quelques exemples :

    • L'ajout, la suppression et la modification des limites de débit sont autorisés :

      Scala
      spark.readStream.format("kafka").option("subscribe", "article")

      à la

      Scala
      spark.readStream.format("kafka").option("subscribe", "article").option("maxOffsetsPerTrigger", ...)

      Pour plus de détails, consultez Configurer la taille de batch Structured Streaming sur Databricks

    • Les modifications apportées aux articles et fichiers abonnés ne sont généralement pas autorisées car les résultats sont imprévisibles : spark.readStream.format("kafka").option("subscribe", "article") à spark.readStream.format("kafka").option("subscribe", "newarticle")

  • Modifications de l’intervalle du trigger : vous pouvez modifier les triggers entre des batches incrémentiels et des intervalles de temps. Voir Modifier les intervalles de trigger entre les exécutions.

  • Modifications dans le type de récepteur de sortie : les modifications entre quelques combinaisons spécifiques de récepteurs sont autorisées. Ceci doit être vérifié au cas par cas. Voici quelques exemples.

    • La cible de fichier vers la cible Kafka est autorisée. Kafka ne verra que les nouvelles données.
    • Le sink Kafka vers le sink de fichiers n'est pas autorisé.
    • Le récepteur Kafka est passé à foreach, et vice versa est autorisé.
  • Modifications des parameters du récepteur de sortie : la possibilité de cette action et la bonne définition de la sémantique de la modification dépendent du récepteur et de la query. Voici quelques exemples.

    • Les modifications apportées au répertoire de sortie d'un sink de fichier ne sont pas autorisées : sdf.writeStream.format("parquet").option("path", "/somePath") vers sdf.writeStream.format("parquet").option("path", "/anotherPath")
    • Les changements de sujet de sortie sont autorisés : sdf.writeStream.format("kafka").option("topic", "topic1") à sdf.writeStream.format("kafka").option("topic", "topic2")
    • Les modifications apportées au récepteur foreach défini par l'utilisateur (c'est-à-dire le code ForeachWriter) sont autorisées, mais la sémantique de la modification dépend du code.
  • Modifications des Opérations de projection/filtre/type map : certains cas sont autorisés. Par exemple :

    • L'ajout/la suppression de filtres est autorisé : de sdf.selectExpr("a") à sdf.where(...).selectExpr("a").filter(...).
    • Les changements dans les projections avec le même schéma de sortie sont autorisés : sdf.selectExpr("stringColumn AS json").writeStream à sdf.select(to_json(...).as("json")).writeStream.
    • Les modifications des projections avec un schéma de sortie différent sont autorisées sous condition : sdf.selectExpr("a").writeStream à sdf.selectExpr("b").writeStream n'est autorisé que si le récepteur de sortie autorise la modification du schéma de "a" à "b".
  • Modifications dans les opérations avec état : Certaines opérations dans les requêtes de streaming doivent maintenir des données d'état afin de mettre à jour continuellement le résultat. Structured Streaming enregistre automatiquement les données d'état dans un stockage tolérant aux pannes (par exemple, DBFS, AWS S3, Azure Blob storage) et les restaure après un redémarrage. Cependant, cela suppose que le schéma des données d'état reste le même après chaque redémarrage. Cela signifie que toute modification (c’est-à-dire les ajouts, les suppressions ou les modifications de schéma) apportée aux opérations avec état d’une query en streaming n’est pas autorisée entre les redémarrages . Voici la liste des opérations avec état dont le schéma ne doit pas être modifié entre les redémarrages afin de garantir la récupération de l’état :

    • Agrégation en streaming : Parsdf.groupBy("a").agg(...) exemple,. Tout changement de nombre ou de type de clés de regroupement ou d'agrégats n'est pas autorisé.
    • Déduplication streaming : Par exemple, sdf.dropDuplicates("a"). Tout changement de nombre ou de type de clés de regroupement ou d'agrégats n'est pas autorisé.
    • Jointure Stream-stream : Par exemple, sdf1.join(sdf2, ...) (c’est-à-dire que les deux entrées sont générées avec sparkSession.readStream). Les modifications du schéma ou des colonnes de jointure d'égalité ne sont pas autorisées. Les modifications du type de jointure (externe ou interne) ne sont pas autorisées. D'autres changements dans la condition de jointure sont mal définis.
    • Opération arbitraire avec état : Par exemple, sdf.groupByKey(...).mapGroupsWithState(...) ou sdf.groupByKey(...).flatMapGroupsWithState(...). Toute modification du schéma de l'état défini par l'utilisateur et du type de délai d'attente n'est pas autorisée. Toute modification au sein de la fonction de mappage d'état définie par l'utilisateur est autorisée, mais l'effet sémantique de la modification dépend de la logique définie par l'utilisateur. Si vous souhaitez réellement prendre en charge les modifications de schéma d'état, vous pouvez explicitement encoder/décoder vos structures de données d'état complexes en octets à l'aide d'un schéma d'encodage/décodage qui prend en charge la migration de schéma. Par exemple, si vous enregistrez votre état sous forme d'octets encodés en Avro, vous pouvez modifier le schéma d'état Avro entre les redémarrages de query, car cela restaure l'état binaire.
important

Les opérateurs avec état dropDuplicates() et dropDuplicatesWithinWatermark() peuvent échouer au redémarrage en raison d'une vérification de compatibilité du schéma d'état lors du changement de mode d'accès au compute.

Le changement entre les modes d'accès dédié et sans isolation est autorisé. Le changement entre les modes d'accès standard et serverless est autorisé. N'essayez pas de changer entre d'autres combinaisons de modes d'accès.

Pour éviter cette erreur, ne modifiez pas le mode d'accès au compute pour les queries de streaming qui contiennent ces opérateurs.

Modifier les sources de streaming avec l'évolution de la source.

Par default, Structured Streaming identifie les sources par leur position dans le plan de query, telles que 0, 1, 2, etc. Toute modification du nombre ou de l’ordre des sources d’entrée interrompt la compatibilité du point de contrôle et nécessite un nouveau point de contrôle. L’évolution de la source vous permet d’attribuer des noms stables et définis par l’utilisateur à chaque source de streaming afin que vous puissiez réorganiser, ajouter ou supprimer des sources d’une query sans perdre l’état du point de contrôle.

L'évolution de la source nécessite Databricks Runtime 18.2 et versions ultérieures.

Configuration requise

Pour activer l'évolution de la source, définissez la configuration Spark suivante :

  • spark.sql.streaming.queryEvolution.enableSourceEvolution: Lorsque true, toutes les sources de streaming dans la query doivent être nommées explicitement à l'aide de l'API .name(). La valeur par default est false.

Définissez la configuration avant de définir la query de streaming :

Python
spark.conf.set("spark.sql.streaming.queryEvolution.enableSourceEvolution", "true")

Règles de nommage

  • Les noms doivent contenir uniquement des caractères alphanumériques et des traits de soulignement ([a-zA-Z0-9_]+).
  • Chaque nom de source doit être unique dans une query.
  • Lorsque l’évolution des sources est activée, chaque source de streaming doit avoir un nom. Des sources sans nom provoquent une UNNAMED_STREAMING_SOURCES_WITH_ENFORCEMENT erreur.

Réorganiser, ajouter et supprimer des sources

Les changements suivants sont sécurisés lors des redémarrages de la query avec le même point de contrôle :

  • Réorganiser les sources : Redémarrez la query avec un ordre de sources différent. Chaque source reprend à partir de son dernier décalage validé en fonction de son nom et ne modifie pas l'état du point de contrôle.
  • **Ajoutez de nouvelles sources** : Redémarrez la query avec une nouvelle source. La nouvelle source traite à partir du début et les sources existantes continuent depuis leurs derniers décalages.
  • Supprimer les sources : redémarrer la query sans la source. La source est définitivement supprimée du point de contrôle. Une source supprimée ne peut pas être ajoutée de nouveau avec le même nom.

Exemple

Utilisez .name() sur DataStreamReader avant d'appeler .load() ou .table():

Python
orders_us = (spark.readStream
.name("orders_us")
.table("catalog.schema.orders_us")
)

orders_eu = (spark.readStream
.name("orders_eu")
.table("catalog.schema.orders_eu")
)

all_orders = orders_us.union(orders_eu)

Limitations

  • Le nommage de la source nécessite un nouveau point de contrôle. Vous ne pouvez pas activer l'évolution de la source sur un point de contrôle existant qui a été créé sans cela.
  • L'évolution de la source est un changement irréversible pour les points de contrôle. Après avoir activé l'évolution de la source pour une query, vous ne pouvez pas désactiver l'évolution de la source et réutiliser le même point de contrôle.
  • Les noms de source sont permanents. Pour renommer une source, supprimez-la, puis ajoutez-la avec un nouveau nom. La source renommée traite depuis le début.