Aller au contenu principal

Créer une application avec état personnalisée

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

remarque

Pour les opérations avec état, telles que les agrégations, la déduplication et les jointures en streaming, Databricks recommande d'utiliser les opérateurs Structured Streaming intégrés au lieu d'une logique personnalisée. Consultez Qu'est-ce que le streaming avec état ?.

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

Exigences

Les opérateurs transformWithState et transformWithStateInPandas ont les exigences suivantes :

  • Disponible dans Databricks Runtime 16.2 et versions ultérieures.

    • Pour le mode temps réel, utilisez Databricks Runtime 17.3 LTS ou une version supé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 default dans Databricks Runtime 17.3 et versions supérieures.

    • Pour Databricks Runtime 17.2 et versions antérieures, vous devez configurer le fournisseur de magasin d'état RocksDB. Databricks recommande d'activer RocksDB dans la configuration Spark.

      Python
      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é à un query Structured Streaming. Vous devez implémenter un processeur avec état personnalisé pour utiliser transformWithState. Structured Streaming inclut des APIs pour créer votre processeur avec état en utilisant Python, Scala ou Java.

Utilisez transformWithState pour appliquer une logique personnalisée à une clé de regroupement. Voici la conception de haut niveau :

  • 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 de la clé sont disponibles en tant qu'itérateur.
  • Utilisez le StatefulProcessorHandle avec des temporisateurs et des conditions définies par l'utilisateur pour contrôler la manière d'émettre des lignes.
  • Pour gérer l’expiration et la taille de l’état, les valeurs d’état prennent en charge des définitions individuelles de durée de vie (TTL).

Puisque transformWithState prend en charge l'évolution des schémas dans le magasin d'état, 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 tenu de retraiter les lignes, ce qui simplifie les déploiements et la maintenance du code. Voir l'évolution des schémas dans le magasin d'état.

important

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

  • PySpark prend en charge à la fois l'API transformWithState basée sur les lignes et l'opérateur transformWithStateInPandas basé sur Pandas.

    • transformWithStateInPandas n'est pas pris en charge en mode temps réel. Utilisez transformWithState à la place. Pour plus de détails, veuillez consulter transformWithState en mode temps réel.
  • Scala ne prend en charge que l'API transformWithState basée sur les lignes.

Les implémentations Scala et Python de transformWithState ont les mêmes capacités, mais avec quelques différences de syntaxe.

Définir un StatefulProcessor

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

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

transformWithState prend en charge trois types d'état : ValueState, ListState et MapState. Chaque type stocke l'état pour chaque clé de regroupement en utilisant 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-batch. Consultez Gérer les lignes d'entrée.
  • Implémentez handleExpiredTimer pour exécuter une logique basée sur le temps, que la clé de regroupement reçoive ou non de nouvelles lignes dans un micro-batch. Voir Gérer les minuteurs expirés.
  • Vous pouvez éventuellement implémenter handleInitialState pour pré-remplir l'état avant que votre application ne 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, mettre, 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 les lignes dans le micro-batch actuel

Oui

Non

Trigger logic basée sur le temps écoulé

Non

Oui

Comportement

handleInputRows

handleExpiredTimer

Obtenir, mettre, 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 les lignes dans le micro-batch actuel

Oui

Non

Trigger logic basée sur le temps écoulé

Non

Oui

Vous pouvez combiner handleInputRows et handleExpiredTimer pour implémenter une logique complexe selon vos besoins.

Par exemple, vous pourriez implémenter une application qui utilise handleInputRows pour mettre à jour les valeurs d'état pour chaque micro-batch et définir un minuteur de 10 secondes dans le futur. Si aucune ligne supplémentaire n’est traitée, vous pouvez utiliser handleExpiredTimer pour émettre les valeurs actuelles du 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 classe StatefulProcessorHandle 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 transmettre le StatefulProcessorHandle à la variable handle. La variable handle lie la variable locale de votre classe Python à la variable d'état.

remarque

Scala utilise la méthode getHandle.

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 en fonction de votre logique d’application complète. Par exemple, vous pourriez suivre les sessions avec un ValueState regroupé 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 carte.

Si votre objet d'état utilise une StructType, vous devez définir des noms uniques pour chaque champ de la structure pour le schéma. Ces noms sont visibles lors de la lecture du magasin d'état. Consultez Lire les informations d'état de Structured Streaming.

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'une struct ou un tuple. Pour ValueState, vous devez implémenter une logique pour remplacer la valeur entière.

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 reset.

ListState

ListState stocke une liste pour chaque clé de regroupement.

Un état de liste est une collection de valeurs, dont chacune peut inclure des types complexes. Chaque valeur dans une liste possède 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 écrasant la liste entière avec un put. Pour reset la durée de vie, vous devez utiliser une opération put.

MapState

MapState stocke un mappage pour chaque clé de regroupement. Les mappages sont l'équivalent Apache Spark d'un dictionnaire Python (dict).

Un état de mappage est une collection de clés distinctes qui sont chacune mappées à une valeur, chacune pouvant inclure des types complexes. Chaque paire clé-valeur d’une carte a sa propre durée de vie.

Vous pouvez mettre à jour la valeur d'une clé spécifique, ou vous pouvez supprimer une clé et sa valeur. Vous pouvez retourner une valeur individuelle en utilisant sa clé, lister toutes les clés, lister toutes les valeurs, ou retourner 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 GROUP BY de la query Structured Streaming. Les états de carte peuvent contenir un nombre arbitraire de paires clé-valeur pour une clé de regroupement.

Par exemple, si votre query utilise GROUP BY user_id et que vous souhaitez définir une carte pour chaque session_id, votre clé de regroupement est user_id et la clé MapState 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

Créez 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 les objets d’état dans votre logique personnalisée. Définissez et initialisez les variables d’état en remplaçant la méthode intégrée init de la classe StatefulProcessor.

Vous pouvez définir un nombre illimité d'objets d'état à 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
  • Un schéma
    • En Python, vous devez spécifier le schéma.
    • En Scala, vous pouvez passer un Encoder pour spécifier le schéma d'état.

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

remarque

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

Utilisez vos variables d'état dans les 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 possède des informations d'état dédiées.

  • Le StatefulProcessor émet des lignes en fonction de votre logique personnalisée et du schéma de sortie spécifié. Voir Émettre des lignes.
  • Utilisez le lecteur statestore pour accéder aux valeurs dans le magasin d'état. Ce lecteur est destiné aux charges de travail batch et n'est pas destiné aux charges de travail à faible latence. Consultez Lire les informations d'état de Structured Streaming.
  • La logique spécifiée à l'aide de handleInputRows s'exécute uniquement si des lignes pour la clé sont présentes dans un micro-batch. Consultez Gérer les lignes d'entrée.
  • Utilisez handleExpiredTimer pour implémenter une logique basée sur le temps qui ne dépend pas de l'observation des lignes pour se déclencher. Voir 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 des lignes associées à une clé de regroupement différente.
  • Vous ne pouvez pas implémenter de logique qui dépend de la comparaison de valeurs ou de la mise à jour de l'état entre les clés de regroupement.

Vous pouvez comparer des valeurs au sein d'une clé de regroupement. Utilisez un MapState pour implémenter une logique avec une seconde clé que votre logique personnalisée peut utiliser. Par exemple, en regroupant par user_id et en utilisant ip_address pour votre clé MapState, vous pouvez suivre les sessions utilisateur simultanées.

Considérations avancées pour l'utilisation de l'état

Les mises à jour de l'état sont tolérantes aux pannes. Si une tâche se bloque avant la fin du traitement d'un micro-batch, la relance utilise la valeur du dernier micro-batch réussi.

Pour des performances optimisées, Databricks vous recommande de traiter toutes les valeurs de l'itérateur pour une clé donnée et de commit les mises à jour en une seule écriture. Lorsque vous écrivez dans une variable d'état, cela Trigger une écriture vers RocksDB.

Les valeurs d'état n'ont pas de valeurs default. Si votre logique nécessite la lecture d'informations d'état existantes, utilisez la méthode exists.

Pour implémenter la logique pour l'état nul, les variables MapState vous permettent de vérifier les clés individuelles ou de lister toutes les clés.

Gérer les lignes d'entrée

Utilisez la méthode handleInputRows pour définir comment votre application traite les lignes et met à jour les valeurs d'état. Cette méthode s'exécute chaque fois que votre query Structured Streaming traite des lignes pour une clé de regroupement.

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

Pour chaque mise à jour de micro-batch traitée, toutes les lignes du micro-batch 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 les valeurs du magasin d'état.

Gérer les temporisateurs expirés

Utilisez la méthode handleExpiredTimer pour implémenter une logique personnalisée basée sur le temps écoulé.

Au sein d'une clé de regroupement, les minuteurs sont identifiés de manière unique par leur timestamp.

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

  • Émission d'informations stockées dans une variable d'état.
  • Éviction des informations d'état stockées.
  • Création d'un nouveau minuteur.

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

Spécifiez le mode temporel

Lorsque vous transmettez votre StatefulProcessor à transformWithState, vous devez spécifier le mode horaire en utilisant le paramètre timeMode.

Les options suivantes sont prises en charge :

Mode Temps

Description

ProcessingTime

Les minuteurs et le TTL sont tous deux pris en charge et sont évalués en fonction de l'heure réelle à laquelle Apache Spark traite chaque micro-batch. Utilisez ProcessingTime lorsque vous souhaitez que les minuteurs se déclenchent à un intervalle fixe par rapport au moment où les lignes sont traitées, indépendamment des Timestamp dans les données.

EventTime

Les minuteurs sont pris en charge et sont évalués en fonction du filigrane d'heure d'événement. Le filigrane avance à mesure qu'Apache Spark observe les Timestamp dans les données d'entrée. Le TTL n'est pas pris en charge avec EventTime. Utilisez EventTime lorsque vos données contiennent des Timestamp et que vous souhaitez que les minuteurs se déclenchent en fonction de la progression de ces Timestamp. Lorsque vous utilisez EventTime, vous devez également spécifier le parameter eventTimeColumnName. Voir eventTimeColumnName.

NoTime OU TimeMode.None()

Les minuteurs et les TTL ne sont pas pris en charge. Utilisez NoTime lorsque votre application à état ne nécessite pas de logique temporelle.

Mode Temps

Description

ProcessingTime

Les minuteurs et le TTL sont tous deux pris en charge et sont évalués en fonction de l'heure réelle à laquelle Apache Spark traite chaque micro-batch. Utilisez ProcessingTime lorsque vous souhaitez que les minuteurs se déclenchent à un intervalle fixe par rapport au moment où les lignes sont traitées, indépendamment des Timestamp dans les données.

EventTime

Les minuteurs sont pris en charge et sont évalués en fonction du filigrane d'heure d'événement. Le filigrane avance à mesure qu'Apache Spark observe les Timestamp dans les données d'entrée. Le TTL n'est pas pris en charge avec EventTime. Utilisez EventTime lorsque vos données contiennent des Timestamp et que vous souhaitez que les minuteurs se déclenchent en fonction de la progression de ces Timestamp. Lorsque vous utilisez EventTime, vous devez également spécifier le parameter eventTimeColumnName. Voir eventTimeColumnName.

NoTime OU TimeMode.None()

Les minuteurs et les TTL ne sont pas pris en charge. Utilisez NoTime lorsque votre application à état ne nécessite pas de logique temporelle.

eventTimeColumnName

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

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

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

Valeurs de minuterie intégrées

Databricks déconseille d'invoquer l'horloge système dans votre application à état personnalisé, car cela peut entraîner des nouvelles tentatives peu fiables en cas d'échec des tâches. Utilisez les méthodes de la classe TimerValues lorsque vous devez accéder au temps de traitement ou au filigrane :

TimerValues

Description

getCurrentProcessingTimeInMs

Renvoie le timestamp de l’heure de traitement du batch actuel en millisecondes depuis l’époque.

getCurrentWatermarkInMs

Renvoie le Timestamp du watermark pour le batch actuel en millisecondes depuis l'époque.

TimerValues

Description

getCurrentProcessingTimeInMs

Renvoie le timestamp de l’heure de traitement du batch actuel en millisecondes depuis l’époque.

getCurrentWatermarkInMs

Renvoie le Timestamp du watermark pour le batch actuel en millisecondes depuis l'époque.

remarque

Le temps de traitement décrit le temps pendant lequel le micro-batch est traité par Apache Spark. De nombreuses sources de streaming, telles que Kafka, incluent également le temps de traitement du système.

Les filigranes sur les requêtes de streaming sont souvent définis par rapport au temps d'événement ou au temps de traitement de la source de streaming. 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 pourriez implémenter des fonctionnalités similaires dans votre application d'état personnalisée en tirant parti de la fonctionnalité TTL, des minuteries et de MapState ou ListState.

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

Pour éviter les erreurs de mémoire insuffisante et pour supprimer les valeurs de type d'état obsolètes, transformWithState prend en charge une valeur de temps de vie (TTL) facultative pour chaque valeur de type d'état. Après expiration, le TTL évince silencieusement les valeurs de type d'état. Le TTL n'exécute pas handleExpiredTimer ou toute logique personnalisée. Pour exécuter le code lorsque l'état expire, utilisez plutôt une minuterie.

important

Si vous n'implémentez pas de TTL, vous devez gérer l'éviction d'état pour éviter les erreurs de mémoire insuffisante.

Pour tous les types d'état, la TTL Reset lors de la mise à jour des informations d'état. Le TTL est appliqué pour chaque valeur de type d'état, avec des règles différentes pour chaque type d'état :

  • Les variables d'état sont limitées aux clés de regroupement.

  • Pour ValueState objets, une seule valeur est stockée par clé de regroupement. Le TTL s'applique à cette valeur.

  • Pour ListState objets, la liste peut contenir de nombreuses valeurs. Le TTL s'applique à chaque valeur d'une liste indépendamment.

    • Bien que la durée de vie (TTL) s'applique à des valeurs individuelles dans un ListState, la seule façon de mettre à jour une valeur individuelle est d'utiliser la méthode put, qui écrase l'intégralité du contenu de la variable ListState et Reset la durée de vie (TTL) pour toutes les valeurs de la liste.
  • Pour les objets MapState, chaque clé de carte a une valeur d'état associée. Le TTL 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 l’éviction d’état, y compris l’émission de lignes. En option, vous pouvez utiliser des minuteurs pour effacer les informations d’état pour une valeur d’état donnée, et émettre des valeurs ou Trigger une logique conditionnelle. Voir Gérer les minuteurs expirés.

Exemple d'application avec état

L'exemple suivant définit un processeur avec état personnalisé, SimpleCounterProcessor, incluant des variables d'état d'exemple. SimpleCounterProcessor utilise ValueState, ListState et MapState pour compter les lignes pour chaque clé de regroupement.

Python
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...
)

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

remarque

En Python, les valeurs d’état sont des tuples. Passez des tuples à put et update, et attendez des tuples de get.

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

Python
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

Utilisez cette approche pour les éléments d'un ListState ou les valeurs d'un MapState également.

Émettre des lignes

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

Les applications personnalisées avec état ne font aucune hypothèse quant à la manière d'utiliser les informations d'état. Pour une condition donnée, l'application peut émettre aucune ligne, une seule ligne ou plusieurs lignes.

remarque

Vous pouvez implémenter plusieurs valeurs d'état et définir plusieurs conditions pour l'émission de lignes, mais toutes les lignes doivent utiliser le même schéma.

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

Émettre des lignes à l'aide d'un objet DataFrame pandas et yield.

Vous pouvez, en option, yield 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 à null.

Gérer l'état initial

Facultativement, vous pouvez transmettre un état initial au premier micro-batch.

Par exemple, vous pourriez l'utiliser pour :

  • Migrer un workflow existant vers une nouvelle application personnalisée.
  • Mettez à niveau un opérateur avec état pour changer votre schéma ou votre logique.
  • Réparez une défaillance qui ne peut pas être réparée automatiquement et qui nécessite une intervention manuelle.
remarque

Utilisez le lecteur de magasin d'état pour query les information d'état à partir d'un point de contrôle existant. Consultez Lire les informations d'état de Structured Streaming.

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

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 distincte StatefulProcessorWithInitialState.

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

Python
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...
)

Utiliser transformWithState dans les LakeFlow Pipelines

Utilisez l'opérateur transformWithState au sein des LakeFlow Pipelines pour implémenter une logique d'état arbitraire dans vos pipelines de streaming en utilisant Python.

Pour ce faire, suivez les étapes suivantes :

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

Pour un exemple qui utilise transformWithState pour surveiller les pulsations des capteurs, consultez Exemple : Utiliser transformWithState pour surveiller les pulsations des capteurs.