Générer une application avec état personnalisé

Vous pouvez utiliser transformWithState pour créer des applications de diffusion en continu avec état et implémenter des solutions à faible latence et en quasi-temps réel. Avec des opérateurs avec état personnalisés, vous pouvez créer une logique avec état arbitraire qui vous permet de créer de nouveaux cas d’usage opérationnels qui ne sont pas possibles avec le traitement structuré traditionnel de streaming.

Remarque

Pour les opérations avec état telles que les agrégations, la déduplication et les jointures de diffusion en continu, Databricks recommande d’utiliser des opérateurs de streaming structuré intégrés au lieu d’une logique personnalisée. Consultez Qu’est-ce que la diffusion en continu avec état ?.

Databricks recommande d’utiliser transformWithState au lieu d’opérateurs hérités, tels que flatMapGroupsWithState et mapGroupsWithState, pour les transformations d’état arbitraires. Consultez Opérateurs avec état arbitraire hérités.

Spécifications

Les opérateurs transformWithState et transformWithStateInPandas répondent aux exigences suivantes :

  • Disponible dans Databricks Runtime 16.2 et versions ultérieures.
    • Pour le mode en temps réel, utilisez Databricks Runtime 17.3 LTS ou version ultérieure. Voir le mode temps réel dans Structured Streaming.
    • Pour le mode d’accès standard, Python est disponible dans Databricks Runtime 16.3 et versions ultérieures, et Scala est disponible dans Databricks Runtime 17.3 et versions ultérieures.
  • RocksDB est le fournisseur de magasin d’état par défaut dans Databricks Runtime 17.3 et versions ultérieures.
    • Pour Databricks Runtime 17.2 et ci-dessous, vous devez configurer le fournisseur de magasin d’état RocksDB. Databricks recommande d’activer RocksDB dans la configuration Spark.

      spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
      

Qu’est-ce que transformWithState ?

L’opérateur transformWithState applique un processeur avec état personnalisé à une requête Structured Streaming. Vous devez implémenter un processeur avec état personnalisé pour utiliser transformWithState. Structured Streaming inclut des API pour la création de votre processeur avec état à l’aide de Python, Scala ou Java.

Permet transformWithState d’appliquer une logique personnalisée à une clé de regroupement. Les éléments suivants décrivent la conception générale :

  • Définissez une ou plusieurs variables d’état.
  • Les informations d’état persistent pour chaque clé de regroupement. Vous pouvez accéder à chaque variable d’état dans le code défini par l’utilisateur.
  • Pour chaque micro-batch traité, toutes les lignes associées à la clé sont disponibles sous forme d’itérateur.
  • Utilisez le StatefulProcessorHandle avec des temporisateurs et des conditions définies par l’utilisateur pour contrôler la manière dont les lignes sont émises.
  • Pour gérer l’expiration de l’état et la taille de l’état, les valeurs d’état prennent en charge les définitions de durée de vie (TTL) individuelles.

Étant donné que transformWithState prend en charge l’évolution du schéma dans le magasin d’états, vous pouvez itérer et mettre à jour vos applications de production sans perdre les informations d’état historiques. Après avoir mis à jour le schéma d’état, vous n’êtes pas obligé de retraiter les lignes, ce qui simplifie les déploiements de code et la maintenance. Consultez Évolution du schéma dans le magasin d’état.

Important

Azure Databricks documentation utilise transformWithState pour décrire les implémentations Python et Scala :

  • PySpark prend en charge l’API basée sur les lignes transformWithState et l’opérateur basé sur Pandas transformWithStateInPandas.
  • Scala prend uniquement en charge l’API basée sur les lignes transformWithState.

Les implémentations en Scala et en Python de transformWithState offrent les mêmes fonctionnalités, mais présentent quelques différences de syntaxe.

Définition d’un StatefulProcessor

Vous définissez un processeur avec état en étendant la StatefulProcessor classe et en implémentant ses méthodes.

Spark transmet un StatefulProcessorHandle à la méthode init de votre StatefulProcessor. Utilisez le descripteur pour créer des variables d’état et interagir avec le magasin d’état.

transformWithState prend en charge trois types d’état : ValueState, ListStateet MapState. Chaque type stocke l’état de chaque clé de regroupement à l’aide d’une structure de données sous-jacente différente.

Implémentez les méthodes suivantes pour définir votre logique personnalisée :

  • Implémentez handleInputRows pour contrôler la façon dont votre application traite les données, met à jour l’état et émet des lignes pour chaque micro-lot. Consultez Gérer les lignes d’entrée.
  • Implémentez handleExpiredTimer pour exécuter une logique temporelle, que la clé de regroupement reçoive ou non de nouvelles lignes dans un micro-lot. Consultez Gérer les minuteurs expirés.
  • Si vous le souhaitez, implémentez handleInitialState pour préremplir l’état avant que votre application traite les lignes d’entrée. Voir Gérer l’état initial.

Le tableau suivant compare les comportements fonctionnels de ces méthodes :

Comportement handleInputRows handleExpiredTimer
Obtenir, placer, mettre à jour ou effacer les valeurs d’état Oui Oui
Créer ou supprimer un minuteur Oui Oui
Émettre des lignes Oui Oui
Itérer sur des lignes dans le micro-batch actuel Oui Non
Logique de déclenchement basée sur le temps écoulé Non Oui

Vous pouvez combiner les deux handleInputRows et handleExpiredTimer implémenter une logique complexe en fonction des besoins.

Par exemple, vous pouvez implémenter une application qui utilise handleInputRows pour mettre à jour les valeurs d’état pour chaque micro-lot et définir un minuteur de 10 secondes à l’avenir. Si aucune ligne supplémentaire n’est traitée, vous pouvez utiliser handleExpiredTimer pour émettre les valeurs actuelles dans le magasin d’état. Si de nouvelles lignes sont traitées pour la clé de regroupement, vous pouvez effacer le minuteur existant et définir un nouveau minuteur.

StatefulProcessorHandle

Dans PySpark, la StatefulProcessorHandle classe vous permet d’accéder aux fonctions qui contrôlent la façon dont votre code utilise les informations d’état.

Lors de l’initialisation d’un StatefulProcessor, vous devez toujours importer et passer le StatefulProcessorHandle à la variable handle. La variable handle lie la variable locale dans votre classe Python à la variable d’état.

Remarque

Scala utilise la getHandle méthode.

Types d’état personnalisés

Vous pouvez implémenter plusieurs objets d’état dans un seul opérateur avec état.

Choisissez un type d’état basé sur votre logique d’application complète. Par exemple, vous pouvez suivre des sessions avec un ValueState, regroupées par user_id et session_id. Ou, pour évaluer les conditions sur plusieurs sessions, utilisez un MapState regroupé par user_id, avec session_id comme clé de mappage.

Si votre objet d’état utilise un StructType, vous devez définir des noms uniques pour chaque champ du struct pour le schéma. Ces noms sont visibles lors de la lecture du magasin d’état. Consultez les informations sur l’état de la diffusion en continu structurée.

Les sections suivantes décrivent les types d’état pris en charge par transformWithState:

ValueState

ValueState stocke une valeur pour chaque clé de regroupement.

Un état de valeur peut inclure des types complexes, tels qu’un struct ou un tuple. Pour ValueState, vous devez implémenter une logique permettant de remplacer l’intégralité de la valeur.

La durée de vie d’un état de valeur est réinitialisée lorsque la valeur est mise à jour. Si vous traitez une clé source pour ValueState sans mettre à jour le ValueState stocké, la durée de vie n’est pas réinitialisée.

ListState

ListState stocke une liste pour chaque clé de regroupement.

Un état de liste est une collection de valeurs, chacune pouvant inclure des types complexes. Chaque valeur d’une liste a sa propre durée de vie.

Vous pouvez ajouter des éléments à une liste en ajoutant des éléments individuels, en ajoutant une liste d’éléments ou en remplaçant toute la liste par un put. Pour réinitialiser la durée de vie, vous devez utiliser une opération put.

MapState

MapState stocke une carte pour chaque clé de regroupement. Les cartes sont l’équivalent d’Apache Spark à un dictionnaire Python (dict).

Un état de mappage est un ensemble de clés distinctes, chacune étant associée à une valeur, laquelle peut inclure des types complexes. Chaque paire clé-valeur dans une carte a sa propre durée de vie.

Vous pouvez mettre à jour la valeur d’une clé spécifique ou supprimer une clé et sa valeur. Vous pouvez retourner une valeur individuelle à l’aide de sa clé, répertorier toutes les clés, répertorier toutes les valeurs ou renvoyer un itérateur pour travailler avec l’ensemble complet de paires clé-valeur dans la carte.

Important

Les clés de regroupement décrivent les champs spécifiés dans la clause de la GROUP BY requête Structured Streaming. Les états de mappage peuvent contenir un nombre arbitraire de paires clé-valeur pour une clé de regroupement.

Par exemple, si votre requête utilise GROUP BY user_id et que vous souhaitez définir une carte pour chacun session_id, votre clé de regroupement est user_id et la MapState clé est session_id:

Python
class SessionTracker(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    self.sessions = handle.getMapState("sessions", StringType(), LongType())

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    for row in rows:
      session_id = row["session_id"]  # session_id is the MapState key
      count = self.sessions.getValue(session_id)[0] if self.sessions.containsKey(session_id) else 0
      new_count = count + 1
      self.sessions.updateValue(session_id, (new_count,))
    yield from []

  def close(self) -> None:
    pass

df.groupBy("user_id").transformWithState(SessionTracker(), ...) # user_id is the grouping key
Langage de programmation Scala
case class Event(userId: String, sessionId: String)

class SessionTracker extends StatefulProcessor[String, Event, Row] {
  @transient private var sessions: MapState[String, Long] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    sessions = getHandle.getMapState[String, Long]("sessions", Encoders.STRING, Encoders.scalaLong, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[Event],
      timerValues: TimerValues): Iterator[Row] = {
    rows.foreach { event =>
      val count = if (sessions.containsKey(event.sessionId)) sessions.getValue(event.sessionId) else 0L
      sessions.updateValue(event.sessionId, count + 1) // sessionId is the MapState key
    }
    Iterator.empty
  }
}

df.as[Event]
  .groupByKey(_.userId) // userId is the grouping key
  .transformWithState(new SessionTracker(), TimeMode.None(), OutputMode.Update())

Créer une variable d’état personnalisée dans le StatefulProcessor

Lorsque vous initialisez votre StatefulProcessor, vous créez une variable locale pour chaque objet d’état qui vous permet d’interagir avec des objets d’état dans votre logique personnalisée. Définissez et initialisez des variables d’état en redéfinissant la méthode intégrée init dans la classe StatefulProcessor.

Vous pouvez définir autant d’objets d’état que vous le souhaitez à l’aide des méthodes getValueState, getListState et getMapState dans votre StatefulProcessor.

Chaque objet d’état doit avoir les éléments suivants :

  • Un nom unique
  • Schéma
    • Dans Python, vous devez spécifier le schéma.
    • En Scala, vous pouvez passer un Encoder pour spécifier le schéma d’état.

Si vous le souhaitez, vous pouvez également fournir une durée de vie (TTL) en millisecondes. Si vous implémentez un état de carte, vous devez fournir une définition de schéma distincte pour les clés de carte et les valeurs.

Remarque

La StatefulProcessor logique gère séparément l’interrogation, la mise à jour et l’émission d’informations d’état. Consultez Utiliser vos variables d’état dans les méthodes avec une logique personnalisée.

Utiliser vos variables d’état dans des méthodes avec une logique personnalisée

Les objets d’état ont des méthodes pour obtenir l’état, mettre à jour les informations d’état existantes et effacer l’état actuel.

Chaque clé de regroupement contient des informations d’état dédiées.

  • Le StatefulProcessor génère des lignes en fonction de votre logique personnalisée et du schéma de sortie spécifié. Voir Émettre des lignes.
  • Utilisez le statestore lecteur pour accéder aux valeurs du stockage d’état. Ce lecteur est destiné aux charges de travail par lots et n’est pas destiné aux charges de travail à faible latence. Consultez les informations sur l’état de la diffusion en continu structurée.
  • La logique spécifiée à l’aide de handleInputRows n’est exécutée que si des lignes correspondant à la clé sont présentes dans un micro-lot. Consultez Gérer les lignes d’entrée.
  • Utilisez handleExpiredTimer pour implémenter une logique temporelle qui ne dépend pas de l’observation des lignes pour se déclencher. Consultez Gérer les minuteurs expirés.

Remarque

Les objets d’état sont isolés en regroupant les clés avec les implications suivantes :

  • Les valeurs d’état ne peuvent pas être affectées par les lignes associées à une autre clé de regroupement.
  • Vous ne pouvez pas implémenter la logique qui dépend de la comparaison des valeurs ou de la mise à jour de l’état entre les clés de regroupement.

Vous pouvez comparer des valeurs dans une clé de regroupement. Utilisez un MapState pour implémenter une logique avec une deuxième clé que votre logique personnalisée peut utiliser. Par exemple, regrouper par user_id et utiliser ip_address pour votre clé MapState vous permet de suivre les sessions simultanées des utilisateurs.

Considérations avancées pour travailler avec l'état

Les mises à jour d’état sont tolérantes aux erreurs. Si une tâche se bloque avant qu’un micro-lot ait terminé le traitement, la nouvelle tentative utilise la valeur du dernier micro-lot réussi.

Pour optimiser les performances, Databricks vous recommande de traiter toutes les valeurs de l’itérateur pour une clé donnée et de valider les mises à jour en écriture unique. Lorsque vous écrivez dans une variable d’état, cela déclenche une écriture dans RocksDB.

Les valeurs d’état n’ont pas de valeurs par défaut. Si votre logique nécessite la lecture des informations d’état existantes, utilisez la exists méthode.

Pour implémenter la logique pour l’état Null, MapState les variables vous permettent de rechercher des clés individuelles ou de répertorier toutes les clés.

Gérer les lignes d’entrée

Utilisez la handleInputRows méthode pour définir la façon dont votre application traite les lignes et met à jour les valeurs d’état. Cette méthode s’exécute chaque fois que votre requête Structured Streaming traite les lignes d’une clé de regroupement.

Pour la plupart des applications avec état implémentées avec transformWithState, la logique principale est définie à l’aide de handleInputRows.

Pour chaque mise à jour de micro-lot traitée, toutes les lignes du micro-lot pour une clé de regroupement donnée sont disponibles à l’aide d’un itérateur. La logique définie par l’utilisateur peut interagir avec toutes les lignes du micro-batch actuel et avec les valeurs figurant dans le magasin d’état.

Gérer les minuteurs expirés

Utilisez la méthode pour implémenter une handleExpiredTimer logique personnalisée en fonction du temps écoulé.

Dans une clé de regroupement, les minuteurs sont identifiés de manière unique par leur horodatage.

Lorsqu’un minuteur expire, le résultat est déterminé par la logique implémentée dans votre application. Les modèles courants sont les suivants :

  • Émission d’informations stockées dans une variable d’état.
  • Évacuez les informations d’état stockées.
  • Création d’un minuteur.

Les minuteurs expirés se déclenchent même si aucune ligne de leur clé associée n’est traitée dans un micro-batch.

Spécifier le mode de temps

Lorsque vous transmettez votre StatefulProcessor à transformWithState, vous devez spécifier le mode temporel à l’aide du paramètre timeMode.

Les options suivantes sont prises en charge :

Mode horaire Descriptif
ProcessingTime Les minuteurs et la TTL sont chacun pris en charge et sont évalués en fonction de la durée écoulée au moment où Apache Spark traite chaque micro-batch. Utilisez ProcessingTime quand vous souhaitez que les minuteurs se déclenchent à un intervalle fixe par rapport au moment où les lignes sont traitées, quels que soient les horodatages dans les données.
EventTime Les minuteurs sont pris en charge et sont évalués en fonction du filigrane de l’heure de l’événement. Le filigrane avance à mesure que Apache Spark observe les horodatages dans les données d’entrée. TTL n’est pas pris en charge avec EventTime. Utilisez EventTime lorsque vos données contiennent des horodatages et que vous souhaitez que les minuteurs se déclenchent en fonction de la progression de ces horodatages. Lors de l’utilisation EventTime, vous devez également spécifier le eventTimeColumnName paramètre. Voir eventTimeColumnName.
NoTime ou TimeMode.None() Les minuteries et le TTL ne sont pas pris en charge. Utilisez NoTime quand votre application avec état ne nécessite pas de logique basée sur le temps.

eventTimeColumnName

Lorsque vous utilisez le EventTime mode de temps, le eventTimeColumnName paramètre spécifie le nom de la colonne dans votre schéma de sortie qui contient l’horodatage d’événement. Apache Spark utilise cette colonne pour propager le watermark vers le flux de sortie, ce qui permet le bon déroulement des opérations en aval basées sur le temps.

Python

eventTimeColumnName est un argument supplémentaire pour transformWithState ou transformWithStateInPandas:

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=MyProcessor(),
      outputStructType=output_schema,
      outputMode="Append",
      timeMode="EventTime",
      eventTimeColumnName="outputTimestamp",
    )
    .writeStream...
)
Langage de programmation Scala

transformWithState accepte eventTimeColumnName à la place de timeMode. Cette approche utilise toujours le EventTime mode :

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(x => x._1)
  .transformWithState(
    new MyProcessor(),
    "outputTimestamp",
    OutputMode.Append(),
  )
  .writeStream...

Valeurs de minuteur intégrées

Databricks recommande de ne pas appeler l’horloge système dans votre application personnalisée avec état, car cela peut entraîner des tentatives de relance peu fiables en cas d’échec de tâches. Utilisez les méthodes de la TimerValues classe lorsque vous devez accéder au temps de traitement ou au filigrane :

TimerValues Descriptif
getCurrentProcessingTimeInMs Retourne l’horodatage de l’heure de traitement du batch actuel en millisecondes depuis l’époque.
getCurrentWatermarkInMs Retourne l’horodatage du filigrane du batch actuel en millisecondes depuis l’époque.

Remarque

Le temps de traitement décrit l’heure à laquelle le micro-lot est traité par Apache Spark. De nombreuses sources de diffusion en continu, telles que Kafka, incluent également le temps de traitement du système.

Les filigranes sur les requêtes de diffusion en continu sont souvent définis par rapport à l’heure de l’événement ou à l’heure de traitement de la source de diffusion en continu. Consultez Appliquer des filigranes pour contrôler les seuils de traitement des données.

Les filigranes et les fenêtres peuvent être utilisés en combinaison avec transformWithState. Vous pouvez implémenter des fonctionnalités similaires dans votre application personnalisée avec état en tirant parti du TTL, des minuteurs et des fonctionnalités MapState ou ListState.

Durée de vie (TTL) pour les types d’état

Pour empêcher les erreurs hors mémoire et supprimer les valeurs de type d’état obsolètes, transformWithState prend en charge une valeur facultative de durée de vie (TTL) pour chaque valeur de type d’état. Après expiration, le TTL évince silencieusement les valeurs du type d’état. TTL n’exécute pas handleExpiredTimer ni aucune logique personnalisée. Pour exécuter du code à l’expiration de l’état, utilisez plutôt un minuteur.

Important

Si vous n’implémentez pas de TTL, vous devez gérer la suppression de l’état pour éviter les erreurs de manque de mémoire.

Pour tous les types d’état, le TTL est réinitialisé lorsqu’on met à jour les informations d’état. Le TTL est appliqué à chaque valeur de type d’état, avec des règles différentes pour chaque type d’état :

  • Les variables d’état sont limitées au regroupement de clés.
  • Pour ValueState les objets, une seule valeur est stockée par clé de regroupement. Le TTL s'applique à cette valeur.
  • Pour ListState les objets, la liste peut contenir de nombreuses valeurs. La durée de vie (TTL) s’applique à chaque valeur d’une liste, indépendamment.
    • Bien que la durée de vie soit étendue à des valeurs individuelles dans un ListState, la seule façon de mettre à jour une valeur individuelle est avec la put méthode, qui remplace l’intégralité du contenu de la ListState variable et réinitialise la durée de vie de toutes les valeurs de la liste.
  • Pour MapState les objets, chaque clé de carte a une valeur d’état associée. La durée de vie s’applique indépendamment à chaque paire clé-valeur dans une carte.

Remarque

Les minuteurs vous permettent de définir une logique personnalisée au-delà de la seule éviction de l’état, notamment pour produire des lignes. Si vous le souhaitez, vous pouvez utiliser des minuteurs pour effacer les informations d’état pour une valeur d’état donnée, et émettre des valeurs ou déclencher une logique conditionnelle. Consultez Gérer les minuteurs expirés.

Exemple d’application avec état

L’exemple suivant définit un processeur avec état personnalisé, SimpleCounterProcessory compris des exemples de variables d’état. SimpleCounterProcessor utilise ValueState, ListStateet MapState pour compter les lignes pour chaque clé de regroupement.

Python (Pandas)

import pandas as pd
from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
    [
        StructField("id", StringType(), True),
        StructField("countAsString", StringType(), True),
    ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    # Schema can also be defined using strings and SQL DDL syntax
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
    count = 0
    for pdf in rows:
      list_state_rows = [(120,), (20,)] # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      pdf_count = pdf.count()
      count += pdf_count.get("value")
    self.value_state.update((count,)) # Count is passed as a tuple
    iter = self.list_state.get()
    list_state_value = next(iter)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,)) # Value is a tuple
    yield pd.DataFrame({"id": key, "countAsString": str(count)})

q = (df.groupBy("key")
  .transformWithStateInPandas(
    statefulProcessor=SimpleCounterProcessor(),
    outputStructType=output_schema,
    outputMode="Update",
    timeMode="None",
  )
  .writeStream...
)

Python (basé sur des lignes)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

output_schema = StructType(
  [
    StructField("id", StringType(), True),
    StructField("countAsString", StringType(), True),
  ]
)

class SimpleCounterProcessor(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    value_state_schema = StructType([StructField("count", IntegerType(), True)])
    list_state_schema = StructType([StructField("count", IntegerType(), True)])
    self.value_state = handle.getValueState(stateName="valueState", schema=value_state_schema)
    self.list_state = handle.getListState(stateName="listState", schema=list_state_schema)
    self.map_state = handle.getMapState(stateName="mapState", userKeySchema="name string", valueSchema="count int")

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    count = 0
    for row in rows:
      list_state_rows = [(120,), (20,)]  # A list of tuples
      self.list_state.put(list_state_rows)
      self.list_state.appendValue((111,))
      self.list_state.appendList(list_state_rows)
      count += 1
    self.value_state.update((count,))  # Count is passed as a tuple
    iter_list = self.list_state.get()
    list_state_value = next(iter_list)[0]
    value = count
    user_key = ("user_key",)
    if self.map_state.exists():
      if self.map_state.containsKey(user_key):
        value += self.map_state.getValue(user_key)[0]
    self.map_state.updateValue(user_key, (value,))  # Value is a tuple
    yield Row(id=key, countAsString=str(count))

q = (
  df.groupBy("key")
    .transformWithState(
      statefulProcessor=SimpleCounterProcessor(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
    )
    .writeStream...
)

Langage de programmation Scala

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.{Dataset, Encoder, Encoders , DataFrame}
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._

spark.conf.set("spark.sql.streaming.stateStore.providerClass","org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")

class SimpleCounterProcessor extends StatefulProcessor[String, (String, String), (String, String)] {
  @transient private var countState: ValueState[Int] = _
  @transient private var listState: ListState[Int] = _
  @transient private var mapState: MapState[String, Int] = _

  private val longEncoder = Encoders.scalaLong
  private val intEncoder = Encoders.scalaInt
  private val stringEncoder = Encoders.STRING

  override def init(
      outputMode: OutputMode,
      timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState",
      intEncoder, TTLConfig.NONE)
    listState = getHandle.getListState[Int]("listState",
      intEncoder, TTLConfig.NONE)
    mapState = getHandle.getMapState[String, Int]("mapState",
      stringEncoder, intEncoder, TTLConfig.NONE)
  }

  override def handleInputRows(
      key: String,
      inputRows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    var count = countState.getOption().getOrElse(0)
    for (row <- inputRows) {
      val listData = Array(120, 20)
      listState.put(listData)
      listState.appendValue(count)
      listState.appendList(listData)
      count += 1
    }
    val iter = listState.get()
    var listStateValue = 0
    if (iter.hasNext) {
      listStateValue = iter.next()
    }
    countState.update(count)
    var value = count
    val userKey = "userKey"
    if (mapState.exists()) {
      if (mapState.containsKey(userKey)) {
        value += mapState.getValue(userKey)
      }
    }
    mapState.updateValue(userKey, value)
    Iterator((key, count.toString))
  }
}

val q = spark
        .readStream
        .format("delta")
        .load("$srcDeltaTableDir")
        .as[(String, String)]
        .groupByKey(x => x._1)
        .transformWithState(
            new SimpleCounterProcessor(),
            TimeMode.None(),
            OutputMode.Update(),
        )
        .writeStream...

Pour plus d’exemples, consultez Exemples d’applications avec état.

Remarque

Dans Python, les valeurs d’état sont des tuples. Transmettez des tuples à put et update, et attendez-vous à recevoir des tuples de get.

Par exemple, si le schéma de votre ValueState est un entier unique :

current_value_tuple = value_state.get() # Returns the value state as a tuple
current_value = current_value_tuple[0]  # Extracts the first item in the tuple
new_value = current_value + 1           # Calculate a new value
value_state.update((new_value,))        # Pass the new value formatted as a tuple

Appliquez également cette approche aux éléments d’un ListState ou aux valeurs d’un MapState.

Émettre des lignes

Vous devez utiliser handleInputRows ou handleExpiredTimer pour définir comment transformWithState émet des lignes pour chaque clé de regroupement. Consultez Gérer les lignes d’entrée et Gérer les minuteurs expirés.

Les applications à état personnalisées ne présupposent rien quant à l’utilisation des informations d’état. Pour une condition donnée, l’application peut émettre aucune ligne, une ligne ou de nombreuses lignes.

Remarque

Vous pouvez implémenter plusieurs valeurs d’état et définir plusieurs conditions pour émettre des lignes, mais toutes les lignes doivent utiliser le même schéma.

Python (Pandas)

Avec transformWithStateInPandas, définissez votre schéma de sortie avec le outputStructType mot clé.

Générer des lignes à l’aide d’un objet pandas DataFrame et yield.

Si vous le souhaitez, vous pouvez yield utiliser un DataFrame vide. Si vous utilisez le mode de sortie update et émettez un DataFrame vide, cela met à jour les valeurs de la clé de regroupement pour qu’elles soient null.

Python (basé sur des lignes)

Avec transformWithState, définissez votre schéma de sortie avec le outputStructType mot clé.

Générez des lignes à l’aide de l’objet Row et de yield.

Si vous le souhaitez, vous pouvez retourner un itérateur vide. Si vous utilisez le mode de sortie update et émettez un itérateur vide, cela met à jour les valeurs de la clé de regroupement de sorte qu’elles soient null.

Langage de programmation Scala

Dans Scala, vous émettez des lignes à l’aide d’un Iterator objet. Le schéma dérive automatiquement du schéma des lignes émises.

Si vous le souhaitez, vous pouvez renvoyer un Iterator vide. Si vous utilisez le mode de sortie update et émettez un Iterator vide, cela met à jour les valeurs de la clé de regroupement pour qu’elles soient null.

Gérer l’état initial

Si vous le souhaitez, vous pouvez passer un état initial au premier micro-lot.

Par exemple, vous pouvez utiliser ceci pour :

  • Migrez un flux de travail existant vers une nouvelle application personnalisée.
  • Mettez à niveau un opérateur avec état pour modifier votre schéma ou votre logique.
  • Réparer une défaillance qui ne peut pas être réparée automatiquement et nécessite une intervention manuelle.

Remarque

Utilisez le lecteur du magasin d’état pour interroger les informations d’état à partir d’un point de contrôle existant. Consultez les informations sur l’état de la diffusion en continu structurée.

Si vous convertissez une table Delta existante en application avec état, lisez la table à l’aide spark.read.table("table_name") et transmettez le DataFrame résultant. Vous pouvez optionnellement sélectionner ou modifier des champs pour les adapter à votre nouvelle application stateful.

Vous fournissez un état initial à l’aide d’un DataFrame avec le même schéma de clé de regroupement que les lignes d’entrée.

Remarque

Python utilise handleInitialState pour spécifier l’état initial lors de la définition d’un StatefulProcessor. Scala utilise la classe StatefulProcessorWithInitialStatedistincte .

L’exemple suivant déclenche un compteur par clé à partir d’une table Delta existante :

Python (basé sur des lignes)

from pyspark.sql import Row
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from typing import Iterator

class CounterWithInitialState(StatefulProcessor):
  def init(self, handle: StatefulProcessorHandle) -> None:
    state_schema = StructType([StructField("count", IntegerType(), True)])
    self.count_state = handle.getValueState("countState", state_schema)

  def handleInitialState(self, key, initialState: Row, timerValues) -> None:
    self.count_state.update((initialState["count"],))

  def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
    count = self.count_state.get()[0] if self.count_state.exists() else 0
    for _ in rows:
      count += 1
    self.count_state.update((count,))
    yield Row(id=key[0], count=count)

  def close(self) -> None:
    pass

output_schema = StructType([
  StructField("id", StringType(), True),
  StructField("count", IntegerType(), True),
])

# Load existing counts as initial state — must use the same grouping key as the input
initial_state = spark.read.table("existing_counts").groupBy("id")

q = (
  df.groupBy("id")
    .transformWithState(
      statefulProcessor=CounterWithInitialState(),
      outputStructType=output_schema,
      outputMode="Update",
      timeMode="None",
      initialState=initial_state,
    )
    .writeStream...
)

Langage de programmation Scala

import org.apache.spark.sql.streaming._
import org.apache.spark.sql.Encoders

class CounterWithInitialState
    extends StatefulProcessorWithInitialState[String, (String, String), (String, String), (String, Int)] {

  @transient private var countState: ValueState[Int] = _

  override def init(outputMode: OutputMode, timeMode: TimeMode): Unit = {
    countState = getHandle.getValueState[Int]("countState", Encoders.scalaInt, TTLConfig.NONE)
  }

  override def handleInitialState(
      key: String, initialState: (String, Int), timerValues: TimerValues): Unit = {
    countState.update(initialState._2)
  }

  override def handleInputRows(
      key: String,
      rows: Iterator[(String, String)],
      timerValues: TimerValues): Iterator[(String, String)] = {
    val count = if (countState.exists()) countState.get() else 0
    val newCount = count + rows.size
    countState.update(newCount)
    Iterator((key, newCount.toString))
  }
}

// Load existing counts as initial state — must use the same grouping key as the input
val initialState = spark.read.table("existing_counts")
  .as[(String, Int)]
  .groupByKey(_._1)

val q = spark
  .readStream
  .format("delta")
  .load(srcDeltaTableDir)
  .as[(String, String)]
  .groupByKey(_._1)
  .transformWithState(
    new CounterWithInitialState(),
    TimeMode.None(),
    OutputMode.Update(),
    initialState,
  )
  .writeStream...

Utiliser transformWithState dans les pipelines Lakeflow

Utilisez l’opérateur transformWithState dans les pipelines Lakeflow pour implémenter une logique arbitraire avec état dans vos pipelines de streaming à l’aide de Python.

Pour ce faire, procédez comme suit :

  1. Définissez le schéma de sortie et la logique de processeur avec état pour vos transformations arbitraires avec état. Pour obtenir des exemples, consultez Exemples d’applications avec état.
  2. Créez un flux de pipeline Lakeflow qui appelle l’opérateur transformWithState sur un DataFrame. Consultez le tutoriel : Créer votre premier pipeline à l’aide de l’éditeur de pipelines Lakeflow.
  3. Exécutez votre pipeline et validez les résultats sur la table ou le récepteur cible.

Pour un exemple qui utilise pour surveiller les battements de capteur, consultez l'exemple : Utiliser pour surveiller les battements de capteur.