Aller au contenu principal

Qu'est-ce que le streaming avec état ?

Cette page explique les query Structured Streaming avec état, y compris les opérations avec état, les recommandations d'optimisation, l'enchaînement de plusieurs opérateurs avec état et le rééquilibrage de l'état.

Une query Structured Streaming *avec état* nécessite des mises à jour incrémentielles des informations d'état intermédiaire, tandis qu'une query Structured Streaming *sans état* suit uniquement les informations relatives aux lignes qui ont été traitées de la source au récepteur. Pour les fonctionnalités d'optimisation disponibles pour les queries sans état, consultez Optimiser les queries de streaming sans état.

Opérations avec état

Les Opérations avec état incluent l'agrégation de streaming, distinct, dropDuplicates, les jointures de Stream à Stream et les applications personnalisées avec état.

Les informations d'état intermédiaires requises pour les requêtes Structured Streaming avec état peuvent entraîner une latence inattendue et des problèmes de production si elles sont mal configurées.

Dans Databricks Runtime 13.3 LTS ou version ultérieure, vous pouvez activer le point de contrôle de changelog avec RocksDB pour réduire la durée du point de contrôle et la latence de bout en bout des charges de travail Structured Streaming. Databricks recommande d'activer le point de contrôle de changelog pour toutes les query Structured Streaming avec état. Voir Activer le point de contrôle de changelog.

Optimiser les requêtes de Structured Streaming avec état

Databricks recommande ce qui suit pour les queries Stateful Structured Streaming :

  • Utilisez des instances optimisées pour le calcul en tant que Worker.
  • Définissez le nombre de partitions de brassage à 1 à 2 fois le nombre de cœurs dans le cluster.
important

Le nombre de partitions de brassage est fixé au moment de la création du point de contrôle. La modification de spark.sql.shuffle.partitions n'a aucun effet sur une query de streaming qui a déjà un point de contrôle — la query continue d'utiliser le nombre de partitions d'origine. Pour appliquer un nouveau nombre de partitions, vous devez start la requête avec un nouvel emplacement de point de contrôle.

Dans Databricks Runtime 18.0 et versions ultérieures, les requêtes de streaming sans état prennent en charge les modifications dynamiques des partitions de shuffle sans nécessiter un nouveau point de contrôle.

Dans Databricks Runtime 18 et versions ultérieures, vous pouvez modifier le nombre de partitions pour les queries avec état sans perdre l'état du checkpoint. Consultez repartitionnement d'état à la demande pour les query de streaming avec état.

  • Définissez la configuration spark.sql.streaming.noDataMicroBatches.enabled sur false dans la SparkSession. Cela empêche le moteur de micro-batch streaming de traiter les micro-batchs qui ne contiennent pas de données. Le fait de définir cette configuration sur false pourrait également entraîner des Opérations avec état qui utilisent des filigranes ou des délais d'attente de traitement ne recevant pas de données de sortie tant que de nouvelles données n'arrivent pas, au lieu de les recevoir immédiatement.

Databricks recommande d'utiliser RocksDB avec le checkpointing du journal des modifications pour gérer l'état des Stream avec état. Voir Configurer le magasin d'état RocksDB sur Databricks.

remarque

Le schéma de gestion d'état ne peut pas être modifié entre les redémarrages de la query. Si une requête a été démarrée avec la gestion default, vous devez la redémarrer à partir de zéro avec un nouvel emplacement de point de contrôle pour modifier le magasin d'état.

Utiliser plusieurs opérateurs avec état dans Structured Streaming

Dans Databricks Runtime 13.3 LTS ou version ultérieure, Databricks offre un support avancé pour les opérateurs stateful dans les charges de travail Structured Streaming. Vous pouvez chaîner plusieurs opérateurs avec état, ce qui signifie que vous pouvez fournir la sortie d'une opération, telle qu'une agrégation fenêtrée, à une autre opération avec état, telle qu'une jointure.

Dans Databricks Runtime 16,2 ou version ultérieure, vous pouvez utiliser transformWithState dans des charges de travail avec plusieurs opérateurs avec état. Consultez Créer une application avec état personnalisée.

Les exemples suivants illustrent plusieurs modèles que vous pouvez utiliser.

important

Les limitations suivantes existent lorsque vous travaillez avec plusieurs opérateurs avec état :

  • Les opérateurs avec état personnalisés hérités (FlatMapGroupWithState et applyInPandasWithState) ne sont pas pris en charge.
  • Seul le mode de sortie d'ajout est pris en charge.

Agrégation chronologique chaînée

Python
words = ...  # streaming DataFrame of schema { timestamp: Timestamp, word: String }

# Group the data by window and word and compute the count of each group
windowedCounts = words.groupBy(
window(words.timestamp, "10 minutes", "5 minutes"),
words.word
).count()

# Group the windowed data by another window and word and compute the count of each group
anotherWindowedCounts = windowedCounts.groupBy(
window(window_time(windowedCounts.window), "1 hour"),
windowedCounts.word
).count()

Agrégation de fenêtre temporelle dans deux Stream différents, suivie d’une jointure de fenêtre Stream à Stream

Python
clicksWindow = clicksWithWatermark.groupBy(
clicksWithWatermark.clickAdId,
window(clicksWithWatermark.clickTime, "1 hour")
).count()

impressionsWindow = impressionsWithWatermark.groupBy(
impressionsWithWatermark.impressionAdId,
window(impressionsWithWatermark.impressionTime, "1 hour")
).count()

clicksWindow.join(impressionsWindow, "window", "inner")

Jointure d’intervalle temporel Stream-Stream suivie d’une agrégation chronologique

Python
joined = impressionsWithWatermark.join(
clicksWithWatermark,
expr("""
clickAdId = impressionAdId AND
clickTime >= impressionTime AND
clickTime <= impressionTime + interval 1 hour
"""),
"leftOuter" # can be "inner", "leftOuter", "rightOuter", "fullOuter", "leftSemi"
)

joined.groupBy(
joined.clickAdId,
window(joined.clickTime, "1 hour")
).count()

Rééquilibrage d'état pour Structured Streaming

Le rééquilibrage de l'état est activé par default pour toutes les charges de travail de streaming dans les LakeFlow Pipelines. Dans Databricks Runtime 11.3 LTS ou version ultérieure, vous pouvez définir l'option de configuration suivante dans la configuration du cluster Spark pour activer le rééquilibrage d'état :

ini
spark.sql.streaming.statefulOperator.stateRebalancing.enabled true

Le rééquilibrage d'état profite aux pipelines Structured Streaming avec état qui subissent des événements de redimensionnement de clusters. Les Opérations de streaming sans état n'en bénéficient pas, quelle que soit la modification de la taille des clusters.

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.

Les événements de redimensionnement de cluster Trigger le rééquilibrage de l'état. Les micro-batchs peuvent avoir une latence plus élevée pendant les événements de rééquilibrage, car l’état est chargé du cloud storage vers les nouveaux exécuteurs.