Aller au contenu principal

Suivi de progression asynchrone

Le suivi asynchrone de la progression réduit la latence pour les pipelines Structured Streaming en permettant aux queries de mettre à jour de manière asynchrone la progression du checkpoint et de traiter les données dans chaque micro-batch.

Lors du traitement des requêtes, Structured Streaming persiste et gère les décalages pour mesurer la progression des requêtes dans le offsetLog et le commitLog dans chaque micro-lot. Sans suivi de progression asynchrone, les Opérations de gestion des décalages affectent directement la latence de traitement, car le traitement des données ne peut pas continuer tant qu'elles ne sont pas terminées.

Suivi asynchrone de la progression

remarque

Le suivi asynchrone de la progression n’est pas compatible avec les Trigger Trigger.once ou Trigger.availableNow. S’il est activé, les queries Structured Streaming avec Trigger.once ou Trigger.availableNow échouent.

Options de configuration

Option

Par défaut

Description

asyncProgressTrackingEnabled

false

Pour activer le suivi asynchrone de la progression.

asyncProgressTrackingCheckpointIntervalMs

1000

L'intervalle en millisecondes entre les écritures pour les décalages et les commits d'achèvement.

Option

Par défaut

Description

asyncProgressTrackingEnabled

false

Pour activer le suivi asynchrone de la progression.

asyncProgressTrackingCheckpointIntervalMs

1000

L'intervalle en millisecondes entre les écritures pour les décalages et les commits d'achèvement.

Activez le suivi asynchrone de la progression

Pour activer le suivi asynchrone des progrès, définissez asyncProgressTrackingEnabled sur true:

Python
stream = (spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "in")
.load()
)

query = (stream.writeStream
.format("kafka")
.option("topic", "out")
.option("checkpointLocation", "/tmp/checkpoint")
.option("asyncProgressTrackingEnabled", "true")
.start()
)

Améliorer le throughput avec la fréquence des points de contrôle.

La fréquence de point de contrôle default de 1000 millisecondes offre un bon throughput pour la plupart des queries. Lorsque les opérations de gestion des décalages se produisent plus rapidement que le suivi asynchrone de la progression ne peut les traiter, un arriéré d'opérations de gestion des décalages se forme. Pour éviter que le backlog ne s'accroisse davantage, le suivi asynchrone des progrès peut bloquer ou ralentir le traitement des données, ce qui pourrait nuire aux avantages attendus en termes de latence.

Dans ce scénario, Databricks vous recommande d’augmenter l’intervalle de point de contrôle :

Python
query = (stream.writeStream
.format("kafka")
.option("topic", "out")
.option("checkpointLocation", "/tmp/checkpoint")
.option("asyncProgressTrackingEnabled", "true")
.option("asyncProgressTrackingCheckpointIntervalMs", "5000")
.start()
)
remarque

Le temps de récupération après une défaillance augmente avec le temps d'intervalle des points de contrôle. En cas de défaillance, un pipeline doit retraiter toutes les données depuis le dernier point de contrôle réussi. Avant d'apporter cette modification en production, tenez compte du compromis entre une latence plus faible pendant le traitement régulier et le temps de récupération en cas de défaillance.

Désactiver le suivi de progression asynchrone

Lorsque le suivi asynchrone de la progression est activé, le Stream ne garantit pas la progression des points de contrôle pour chaque batch. Vous devez créer un point de contrôle de la progression avant de pouvoir désactiver cette fonctionnalité.

Pour désactiver, suivez ces étapes :

  1. Traiter au moins deux micro-batchs avec asyncProgressTrackingEnabled défini sur true et asyncProgressTrackingCheckpointIntervalMs défini sur 0:
Python
query = (stream.writeStream
.format("kafka")
.option("topic", "out")
.option("checkpointLocation", "/tmp/checkpoint")
.option("asyncProgressTrackingEnabled", "true")
.option("asyncProgressTrackingCheckpointIntervalMs", "0")
.start()
)
  1. Arrêter la query :
Python
query.stop()
  1. Désactivez le suivi asynchrone de la progression et redémarrez la query :
Python
query = (stream.writeStream
.format("kafka")
.option("topic", "out")
.option("checkpointLocation", "/tmp/checkpoint")
.option("asyncProgressTrackingEnabled", "false")
.start()
)

Si vous désactivez le suivi de progression asynchrone sans suivre les étapes ci-dessus, vous pourriez rencontrer l'erreur suivante :

java.lang.IllegalStateException: batch x doesn't exist

Dans les Logs du Driver, vous pourriez voir l'erreur suivante :

The offset log for batch x doesn't exist, which is required to restart the query from the latest batch x from the offset log. Please ensure there are two subsequent offset logs available for the latest batch via manually deleting the offset file(s). Please also ensure the latest batch for commit log is equal or one batch earlier than the latest batch for offset log.

Limitations

  • Pour les sinks Kafka, le suivi asynchrone de la progression ne prend en charge que les pipelines stateless.
  • Le suivi asynchrone de la progression ne garantit pas un traitement de bout en bout exactement une seule fois, car les plages de décalage d'un batch peuvent changer en cas d'échec. Certaines destinations, comme Kafka, ne fournissent jamais de garanties d'exécution unique.