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.
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.
Pythonspark.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
StatefulProcessorHandleavec 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.
La documentation Databricks utilise transformWithState pour décrire les implémentations Python et Scala :
-
PySpark prend en charge à la fois l'API
transformWithStatebasée sur les lignes et l'opérateurtransformWithStateInPandasbasé sur Pandas.transformWithStateInPandasn'est pas pris en charge en mode temps réel. UtiliseztransformWithStateà la place. Pour plus de détails, veuillez consultertransformWithStateen mode temps réel.
-
Scala ne prend en charge que l'API
transformWithStatebasé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
handleInputRowspour 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
handleExpiredTimerpour 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
handleInitialStatepour 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 |
|
|
|---|---|---|
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.
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.
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
- Scala
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
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é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
Encoderpour 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.
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
statestorepour 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
handleInputRowss'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
handleExpiredTimerpour 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.
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 |
|---|---|
| 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 |
| 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 |
| Les minuteurs et les TTL ne sont pas pris en charge. Utilisez |
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.
- Python
- Scala
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...
)
transformWithState accepte eventTimeColumnName à la place de timeMode. Cette approche utilise toujours le mode EventTime :
val q = spark
.readStream
.format("delta")
.load(srcDeltaTableDir)
.as[(String, String)]
.groupByKey(x => x._1)
.transformWithState(
new MyProcessor(),
"outputTimestamp",
OutputMode.Append(),
)
.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 :
| Description |
|---|---|
| Renvoie le timestamp de l’heure de traitement du batch actuel en millisecondes depuis l’époque. |
| Renvoie le Timestamp du watermark pour le batch actuel en millisecondes depuis l'époque. |
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.
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
ValueStateobjets, une seule valeur est stockée par clé de regroupement. Le TTL s'applique à cette valeur. -
Pour
ListStateobjets, 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éthodeput, qui écrase l'intégralité du contenu de la variableListStateet Reset la durée de vie (TTL) pour toutes les valeurs de la liste.
- Bien que la durée de vie (TTL) s'applique à des valeurs individuelles dans un
-
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.
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 (Pandas)
- Python (row-based)
- Scala
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...
)
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...
)
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.
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 :
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.
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.
- Python (Pandas)
- Python (row-based)
- Scala
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.
Avec transformWithState, définissez votre schéma de sortie avec le mot-clé outputStructType.
Émettez des lignes à l’aide d’un objet Row et yield.
En option, vous pouvez renvoyer un itérateur vide. Si vous utilisez le mode de sortie update et émettez un itérateur vide, cela met à jour les valeurs pour que la clé de regroupement soit null.
En Scala, vous émettez des lignes à l'aide d'un objet Iterator. Le schéma se dérive automatiquement du schéma des lignes émises.
Vous pouvez éventuellement retourner 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 à 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.
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.
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 (row-based)
- Scala
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...
)
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 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 :
- 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.
- Créez un flux de Lakeflow pipeline qui invoque l'opérateur
transformWithStatesur un DataFrame. Voir Didacticiel : créer votre premier Lakeflow Pipelines à l'aide de l'éditeur de pipelines Lakeflow. - 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.