Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
U kunt transformWithState gebruiken om stateful streamingtoepassingen op te bouwen en om oplossingen met lage latentie en bijna real-time te implementeren. Met aangepaste stateful operators kunt u willekeurige stateful logica maken waarmee u nieuwe operationele use cases kunt bouwen die niet mogelijk zijn met traditionele Structured Streaming-verwerking.
Notitie
Voor statusafhankelijke bewerkingen, zoals aggregaties, deduplicatie en streaming-joins, raadt Databricks aan ingebouwde Structured Streaming-operators te gebruiken in plaats van aangepaste logica. Zie Wat is 'stateful streaming'?.
Databricks raadt aan om in plaats van verouderde operators te gebruiken transformWithState , zoals flatMapGroupsWithState en mapGroupsWithStatevoor willekeurige statustransformaties. Zie Verouderde willekeurige stateful operators.
Eisen
De transformWithState en transformWithStateInPandas operators hebben de volgende vereisten:
- Beschikbaar in Databricks Runtime 16.2 en hoger.
- Voor realtime-modus gebruikt u Databricks Runtime 17.3 LTS of hoger. Bekijk de real-time modus in Structured Streaming.
- Voor de standaardtoegangsmodus is Python beschikbaar in Databricks Runtime 16.3 en hoger en Scala is beschikbaar in Databricks Runtime 17.3 en hoger.
- RocksDB is de standaardprovider voor statusopslag in Databricks Runtime 17.3 en hoger.
Voor Databricks Runtime 17.2 en lager moet u de provider van het rocksDB-statusarchief configureren. Databricks raadt het inschakelen van RocksDB aan in de Spark-configuratie.
spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
Wat is transformWithState?
De operator transformWithState past een aangepaste stateful processor toe op een Structured Streaming-query. U moet een aangepaste stateful processor implementeren om transformWithStatete gebruiken. Structured Streaming bevat API's voor het bouwen van uw stateful processor met behulp van Python, Scala of Java.
Gebruik transformWithState dit om aangepaste logica toe te passen op een groeperingssleutel. Hieronder wordt het ontwerp op hoog niveau beschreven:
- Definieer een of meer statusvariabelen.
- Statusinformatie blijft behouden voor elke groeperingssleutel. U kunt elke statusvariabele openen in door de gebruiker gedefinieerde code.
- Voor elke verwerkte microbatch zijn alle rijen voor de sleutel beschikbaar in de vorm van een iterator.
- Gebruik de
StatefulProcessorHandlefunctie met timers en door de gebruiker gedefinieerde voorwaarden om te bepalen hoe u rijen verzendt. - Als u statusverlooptijd en statusgrootte wilt beheren, ondersteunen statuswaarden afzonderlijke TTL-definities (Time-to-Live).
Omdat transformWithState de ontwikkeling van schema's in het statusarchief wordt ondersteund, kunt u uw productietoepassingen herhalen en bijwerken zonder historische statusgegevens te verliezen. Na het bijwerken van het statusschema hoeft u geen rijen opnieuw te verwerken, wat code-implementaties en onderhoud vereenvoudigt. Bekijk de ontwikkeling van schema's in het statusarchief.
Belangrijk
Azure Databricks documentatie gebruikt transformWithState om zowel Python als Scala-implementaties te beschrijven:
- PySpark ondersteunt zowel de op rijen gebaseerde
transformWithStateAPI als de pandas-operatortransformWithStateInPandas.-
transformWithStateInPandaswordt niet ondersteund in realtimemodus. Gebruik in plaats daarvantransformWithState. Zie de realtime-modus voor meer informatietransformWithState.
-
- Scala ondersteunt alleen de API op basis van
transformWithStaterijen.
De Scala- en Python-implementaties hebben transformWithState dezelfde mogelijkheden, maar met enkele verschillen in syntaxis.
Een StatefulProcessor definiëren
U definieert een stateful processor door de StatefulProcessor klasse uit te breiden en de methoden ervan te implementeren.
Spark geeft een StatefulProcessorHandle door aan de methode init van uw StatefulProcessor. Gebruik de handle om statusvariabelen te maken en met de statusopslag te interageren.
transformWithState ondersteunt drie statustypen: ValueState, ListStateen MapState. Elk type slaat de status voor elke groeperingssleutel op met behulp van een andere onderliggende gegevensstructuur.
Implementeer de volgende methoden om uw aangepaste logica te definiëren:
- Implementeer
handleInputRowsom te bepalen hoe uw toepassing gegevens verwerkt, status bijwerken en rijen verzendt voor elke microbatch. Zie Invoerrijenverwerken. - Implementeer
handleExpiredTimerom logica op basis van tijd uit te voeren, ongeacht of de groeperingssleutel nieuwe rijen in een microbatch ontvangt. Zie Verlopen timers verwerken. - U kunt eventueel implementeren
handleInitialStateom de status vooraf in te vullen voordat uw toepassing invoerrijen verwerkt. Zie Beginstatus Afhandelen.
In de volgende tabel worden de functionele gedragingen van deze methoden vergeleken:
| Gedrag | handleInputRows |
handleExpiredTimer |
|---|---|---|
| Toestandswaarden ophalen, plaatsen, bijwerken of wissen | Ja | Ja |
| Een timer maken of verwijderen | Ja | Ja |
| Rijen verzenden | Ja | Ja |
| Itereren over de rijen in de huidige microbatch | Ja | Nee |
| Triggerlogica op basis van verstreken tijd | Nee | Ja |
U kunt zowel handleInputRows als handleExpiredTimer combineren om indien nodig complexe logica te implementeren.
U kunt bijvoorbeeld een toepassing implementeren die gebruikmaakt van handleInputRows om statuswaarden voor elke microbatch bij te werken en in de toekomst een timer van 10 seconden in te stellen. Als er geen extra rijen worden verwerkt, kunt u handleExpiredTimer gebruiken om de huidige waarden in de statusopslag uit te sturen. Als nieuwe rijen worden verwerkt voor de groeperingssleutel, kunt u de bestaande timer wissen en een nieuwe timer instellen.
StatefulProcessorHandle
In PySpark kunt u met de StatefulProcessorHandle klasse toegang krijgen tot functies die bepalen hoe uw code statusgegevens gebruikt.
Bij het initialiseren van een StatefulProcessor, moet u altijd de StatefulProcessorHandle variabele importeren en doorgeven aan de handle variabele. De variabele handle koppelt de lokale variabele in uw Python klasse aan de statusvariabele.
Notitie
Scala maakt gebruik van de methode getHandle.
aangepaste statustypen
U kunt meerdere statusobjecten implementeren in één stateful operator.
Kies een statustype op basis van de volledige toepassingslogica. U kunt bijvoorbeeld sessies bijhouden met een ValueState, gegroepeerd op basis van user_id en session_id. Of, om voorwaarden over meerdere sessies te evalueren, gebruikt u een MapState gegroepeerd op basis van user_id met session_id als mapsleutel.
Als uw statusobject een StructTypeobject gebruikt, moet u unieke namen definiëren voor elk veld in de struct voor het schema. Deze namen zijn zichtbaar bij het lezen van de statusopslag. Zie Lees informatie over de status van Gestructureerd Streamen.
In de volgende secties worden de statustypen beschreven die worden ondersteund door transformWithState:
ValueState
ValueState slaat een waarde op voor elke groeperingssleutel.
Een waardestatus kan complexe typen bevatten, zoals een struct of tuple. Voor ValueState moet u logica implementeren om de volledige waarde te vervangen.
De time-to-live voor een waardestatus wordt opnieuw ingesteld wanneer de waarde wordt bijgewerkt. Als u een bronsleutel voor ValueState verwerkt zonder de opgeslagen ValueState bij te werken, wordt de time-to-live niet opnieuw ingesteld.
ListState
ListState slaat een lijst op voor elke groeperingssleutel.
Een lijststatus is een verzameling waarden, die elk complexe typen kunnen bevatten. Elke waarde in een lijst heeft een eigen time-to-live.
U kunt items toevoegen aan een lijst door afzonderlijke items toe te voegen, een lijst met items toe te voegen of de hele lijst te overschrijven met een put. Als u time-to-live opnieuw wilt instellen, moet u een put bewerking gebruiken.
MapState
MapState slaat een kaart op voor elke groeperingssleutel. Kaarten zijn het Apache Spark-equivalent aan een Python woordenlijst (dict).
Een mapstatus is een verzameling unieke sleutels die elk aan een waarde zijn gekoppeld, waarvan elk complexe typen kan bevatten. Elk sleutel-waardepaar in een map heeft zijn eigen time-to-live.
U kunt de waarde van een specifieke sleutel bijwerken of u kunt een sleutel en de bijbehorende waarde verwijderen. U kunt een afzonderlijke waarde retourneren met behulp van de sleutel, alle sleutels weergeven, alle waarden weergeven of een iterator retourneren om te werken met de volledige set sleutel-waardeparen in de kaart.
Belangrijk
Met groeperingssleutels worden de velden beschreven die zijn opgegeven in de GROUP BY-component van de Structured Streaming-query. Map-statussen kunnen een willekeurig aantal sleutel-waardeparen bevatten bij een groeperingssleutel.
Als uw query bijvoorbeeld gebruikmaakt van GROUP BY user_id en u voor elke session_id een toewijzing wilt definiëren, is uw groeperingssleutel user_id en is de MapState-sleutel 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
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())
Een aangepaste statusvariabele maken in de StatefulProcessor
Wanneer u uw StatefulProcessorinitialiseert, maakt u een lokale variabele voor elk statusobject waarmee u kunt communiceren met statusobjecten in uw aangepaste logica. Definieer en initialiseer statusvariabelen door de ingebouwde methode in init de StatefulProcessor klasse te overschrijven.
U kunt een willekeurig aantal statusobjecten definiëren met behulp van de getValueState, getListStateen getMapState methoden in uw StatefulProcessor.
Elk statusobject moet het volgende hebben:
- Een unieke naam
- Een schema
- In Python moet u het schema opgeven.
- In Scala kunt u een
Encoderdoorgeven om het statusschema op te geven.
U kunt desgewenst ook een TTL-duur (Time to Live) opgeven in milliseconden. Als u een kaartstatus implementeert, moet u een afzonderlijke schemadefinitie opgeven voor de kaartsleutels en de waarden.
Notitie
De StatefulProcessor logica wordt afzonderlijk verwerkt voor het opvragen, bijwerken en verzenden van statusinformatie. Zie Uw statusvariabelen gebruiken in methoden met aangepaste logica.
Uw statusvariabelen gebruiken in methoden met aangepaste logica
Statusobjecten hebben methoden voor het ophalen van status, het bijwerken van bestaande statusgegevens en het wissen van de huidige status.
Elke groeperingssleutel heeft toegewezen statusinformatie.
- De
StatefulProcessorgenereert rijen op basis van je aangepaste logica en het opgegeven uitvoerschema. Zie Rijen verzenden. - Gebruik de
statestorelezer om toegang te krijgen tot waarden in het statusarchief. Deze lezer is bedoeld voor batchworkloads en is niet bedoeld voor workloads met lage latentie. Zie Lees informatie over de status van Gestructureerd Streamen. - Logica die is opgegeven met
handleInputRowswordt alleen uitgevoerd als rijen voor de sleutel aanwezig zijn in een micro-batch. Zie Invoerrijenverwerken. - Gebruik
handleExpiredTimerdit om op tijd gebaseerde logica te implementeren die niet afhankelijk is van het observeren van rijen die moeten worden geactiveerd. Zie Verlopen timers verwerken.
Notitie
Statusobjecten worden geïsoleerd door sleutels te groeperen met de volgende implicaties:
- Statuswaarden kunnen niet worden beïnvloed door rijen die zijn gekoppeld aan een andere groeperingssleutel.
- Je kan geen logica implementeren die afhankelijk is van het vergelijken van waarden of het bijwerken van de status tussen groeperingssleutels.
U kunt waarden in een groeperingssleutel vergelijken. Gebruik een MapState om logica te implementeren met een tweede sleutel die uw aangepaste logica kan gebruiken. Als u bijvoorbeeld op basis van user_id groepeert en ip_address gebruikt voor uw MapState-sleutel, kunt u gelijktijdige gebruikerssessies bijhouden.
Geavanceerde overwegingen voor het omgaan met status
Statusupdates zijn fouttolerant. Als een taak vastloopt voordat een microbatch is verwerkt, gebruikt de nieuwe poging de waarde van de laatste geslaagde microbatch.
Voor geoptimaliseerde prestaties raadt Databricks u aan om alle waarden in de iterator te verwerken voor een bepaalde sleutel en updates door te voeren in één schrijfbewerking. Wanneer u naar een statusvariabele schrijft, wordt hiermee een schrijfbewerking naar RocksDB geactiveerd.
Statuswaarden hebben geen standaardwaarden. Als voor uw logica bestaande statusgegevens moeten worden gelezen, gebruikt u de exists methode.
Als u logica voor null-status wilt implementeren, MapState kunt u met variabelen controleren op afzonderlijke sleutels of alle sleutels weergeven.
Invoerrijen verwerken
Gebruik de handleInputRows methode om te definiëren hoe uw toepassing rijen verwerkt en statuswaarden bijwerken. Deze methode wordt telkens uitgevoerd wanneer uw Structured Streaming-query rijen verwerkt voor een groeperingssleutel.
Voor de meeste stateful toepassingen die zijn geïmplementeerd met transformWithState, wordt de kernlogica gedefinieerd met behulp van handleInputRows.
Voor elke microbatch-update die wordt verwerkt, zijn alle rijen in de microbatch voor een bepaalde groeperingssleutel beschikbaar met behulp van een iterator. Door de gebruiker gedefinieerde logica kan communiceren met alle rijen uit de huidige microbatch en waarden in de statestore.
Verlopen timers verwerken
Gebruik de handleExpiredTimer methode om aangepaste logica te implementeren op basis van verstreken tijd.
Binnen een groeperingssleutel worden timers uniek geïdentificeerd door hun tijdstempel.
Wanneer een timer verloopt, wordt het resultaat bepaald door de logica die in uw toepassing is geïmplementeerd. Veelvoorkomende patronen zijn:
- Informatie verzenden die is opgeslagen in een statusvariabele.
- Opgeslagen statusgegevens verwijderen.
- Een nieuwe timer maken.
Verlopen timers gaan af, zelfs als er in een micro-batch geen rijen voor de bijbehorende sleutel worden verwerkt.
De tijdmodus opgeven
Wanneer u uw StatefulProcessor aan transformWithState doorgeeft, moet u de tijdmodus opgeven met de parameter timeMode.
De volgende opties worden ondersteund:
| Tijdmodus | Beschrijving |
|---|---|
ProcessingTime |
Timers en TTL worden beide ondersteund en geëvalueerd op basis van de wandkloktijd wanneer Apache Spark elke microbatch verwerkt. Gebruik ProcessingTime deze optie als u wilt dat timers met een vast interval worden geactiveerd ten opzichte van wanneer rijen worden verwerkt, ongeacht tijdstempels in de gegevens. |
EventTime |
Timers worden ondersteund en geëvalueerd op basis van het event-time-watermerk. Het watermerk gaat verder naarmate Apache Spark tijdstempels in de invoergegevens bekijkt. TTL wordt niet ondersteund met EventTime. Gebruik EventTime wanneer uw gegevens tijdstempels bevatten en u wilt dat timers worden geactiveerd op basis van de voortgang van deze tijdstempels. Wanneer u deze EventTimegebruikt, moet u ook de eventTimeColumnName parameter opgeven. Zie eventTimeColumnName. |
NoTime of TimeMode.None() |
Timers en TTL worden niet ondersteund. Gebruik NoTime wanneer uw statusvolle toepassing geen tijdsgebaseerde logica vereist. |
eventTimeColumnName
Wanneer u de EventTime tijdmodus gebruikt, geeft de eventTimeColumnName parameter de naam van de kolom in het uitvoerschema op die de tijdstempel van de gebeurtenis bevat. Apache Spark gebruikt deze kolom om het watermerk door te geven aan de uitvoerstroom, waardoor juiste downstreambewerkingen op basis van tijd worden ingeschakeld.
Python
eventTimeColumnName is een aanvullend argument voor transformWithState of transformWithStateInPandas:
q = (
df.groupBy("key")
.transformWithState(
statefulProcessor=MyProcessor(),
outputStructType=output_schema,
outputMode="Append",
timeMode="EventTime",
eventTimeColumnName="outputTimestamp",
)
.writeStream...
)
Scala
transformWithState accepteert eventTimeColumnName in plaats van timeMode. Deze aanpak gebruikt altijd de modus EventTime:
val q = spark
.readStream
.format("delta")
.load(srcDeltaTableDir)
.as[(String, String)]
.groupByKey(x => x._1)
.transformWithState(
new MyProcessor(),
"outputTimestamp",
OutputMode.Append(),
)
.writeStream...
Ingebouwde timerwaarden
Databricks raadt af om de systeemklok in uw aangepaste stateful toepassing aan te roepen, omdat dit kan leiden tot onbetrouwbare herhalingen bij het mislukken van een taak. Gebruik de methoden in de klasse TimerValues wanneer u toegang moet hebben tot de verwerkingstijd of het watermerk:
TimerValues |
Beschrijving |
|---|---|
getCurrentProcessingTimeInMs |
Retourneert de tijdstempel van de verwerkingstijd voor de huidige batch in milliseconden sinds epoch. |
getCurrentWatermarkInMs |
Retourneert de tijdstempel van het watermerk voor de huidige batch in milliseconden sinds epoch. |
Notitie
De verwerkingstijd beschrijft de tijd die de microbatch door Apache Spark wordt verwerkt. Veel streamingbronnen, zoals Kafka, bevatten ook systeemverwerkingstijd.
Watermerken voor streamingquery's worden vaak gedefinieerd op basis van gebeurtenistijd of de verwerkingstijd van de streamingbron. Zie Watermerken toepassen om drempelwaarden voor gegevensverwerking te beheren.
Zowel watermerken als vensters kunnen worden gebruikt in combinatie met transformWithState. U kunt vergelijkbare functionaliteit implementeren in uw aangepaste stateful toepassing door gebruik te maken van TTL, timers en MapState of ListState functionaliteit.
Time-to-live (TTL) voor statustypen
Als u fouten met onvoldoende geheugen wilt voorkomen en verouderde statustypewaarden wilt verwijderen, transformWithState ondersteunt u een optionele TTL-waarde (Time to Live) voor elke statustypewaarde. Na het verlopen verwijdert TTL stilzwijgend waarden van het statetype. TTL voert handleExpiredTimer of aangepaste logica niet uit. Als u code wilt uitvoeren wanneer de status verloopt, gebruikt u in plaats daarvan een timer.
Belangrijk
Als u TTL niet implementeert, moet u statusverwijdering afhandelen om geheugenfouten te voorkomen.
Voor alle statustypen wordt TTL opnieuw ingesteld bij het bijwerken van statusgegevens. TTL wordt afgedwongen voor elke statustypewaarde, met verschillende regels voor elk statustype:
- Statusvariabelen zijn gericht op het groeperen van sleutels.
- Voor
ValueStateobjecten wordt slechts één waarde per groeperingssleutel opgeslagen. TTL is van toepassing op deze waarde. - Voor
ListStateobjecten kan de lijst veel waarden bevatten. TTL is onafhankelijk van toepassing op elke waarde in een lijst.- Hoewel TTL is gericht op afzonderlijke waarden in een
ListState, is de enige manier om een afzonderlijke waarde bij te werken met deputmethode, die de volledige inhoud van deListStatevariabele overschrijft en TTL opnieuw instelt voor alle waarden in de lijst.
- Hoewel TTL is gericht op afzonderlijke waarden in een
- Voor
MapStateobjecten heeft elke kaartsleutel een bijbehorende statuswaarde. TTL is onafhankelijk van toepassing op elk sleutel-waardepaar in een kaart.
Notitie
Met timers kunt u aangepaste logica definiëren buiten verwijdering van statussen, inclusief het verzenden van rijen. U kunt timers ook gebruiken om statusinformatie voor een bepaalde statuswaarde te wissen en waarden te verzenden of voorwaardelijke logica te activeren. Zie Verlopen timers verwerken.
Voorbeeld van een "stateful" toepassing
In het volgende voorbeeld wordt een aangepaste stateful processor gedefinieerd, SimpleCounterProcessorinclusief voorbeeldstatusvariabelen.
SimpleCounterProcessor gebruikt ValueState, ListStateen MapState om rijen te tellen voor elke groeperingssleutel.
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 (op regels gebaseerd)
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...
)
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...
Zie Voorbeeld van stateful toepassingenvoor meer voorbeelden.
Notitie
In Python zijn statuswaarden tuples. Geef tuples door aan put en update, en verwacht tuples uit get.
Als bijvoorbeeld het schema voor uw ValueState uit één geheel getal bestaat:
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
Gebruik deze methode ook voor items in een ListState of waarden in een MapState .
Rijen verzenden
U moet handleInputRows of handleExpiredTimer gebruiken om te definiëren hoe transformWithState rijen produceert voor elke groeperingssleutel. Zie Invoerrijen verwerken en Verlopen timers verwerken.
Aangepaste toepassingen met status gaan niet uit van hoe statusinformatie wordt gebruikt. Voor een bepaalde voorwaarde kan de toepassing geen rijen, één rij of veel rijen verzenden.
Notitie
U kunt meerdere statuswaarden implementeren en meerdere voorwaarden definiëren voor het verzenden van rijen, maar alle rijen moeten hetzelfde schema gebruiken.
Python (Pandas)
Met transformWithStateInPandas, definieer je je uitvoerschema met het outputStructType sleutelwoord.
Rijen verzenden met behulp van een Pandas DataFrame-object en yield.
U kunt yield eventueel een leeg DataFrame gebruiken. Als u de uitvoermodus gebruikt update en een leeg DataFrame verzendt, worden de waarden voor de groeperingssleutel bijgewerkt null.
Python (op regels gebaseerd)
Met transformWithState, definieer je je uitvoerschema met het outputStructType sleutelwoord.
Genereer rijen met een Row-object en yield.
U kunt eventueel een lege iterator retourneren. Als u de uitvoermodus gebruikt update en een lege iterator verzendt, worden de waarden voor de groeperingssleutel bijgewerkt null.
Scala
In Scala verzendt u rijen met behulp van een Iterator object. Het schema wordt automatisch afgeleid van het schema van de verzonden rijen.
U kunt eventueel een lege Iteratorwaarde retourneren. Als u de update-uitvoermodus gebruikt en een lege Iterator genereert, worden de waarden voor de groeperingssleutel bijgewerkt naar null.
Begintoestand verwerken
U kunt desgewenst een initiële status doorgeven aan de eerste microbatch.
U kunt dit bijvoorbeeld gebruiken om het volgende te doen:
- Een bestaande werkstroom migreren naar een nieuwe aangepaste toepassing.
- Werk een stateful operator bij om uw schema of logica te wijzigen.
- Herstel een fout die niet automatisch kan worden hersteld en waarvoor handmatige tussenkomst is vereist.
Notitie
Gebruik de lezer van het statusarchief om statusgegevens van een bestaand controlepunt op te vragen. Zie Lees informatie over de status van Gestructureerd Streamen.
Als u een bestaande Delta-tabel converteert naar een stateful applicatie, gebruik dan spark.read.table("table_name") om de tabel te lezen en geef het resulterende DataFrame door. U kunt desgewenst velden selecteren of wijzigen om aan te passen aan uw nieuwe statusgevoelige toepassing.
U geeft een initiële status op met behulp van een DataFrame met hetzelfde groeperingssleutelschema als de invoerrijen.
Notitie
Python gebruikt handleInitialState om de initiële status op te geven tijdens het definiëren van een StatefulProcessor. Scala maakt gebruik van de afzonderlijke klasse StatefulProcessorWithInitialState.
In het volgende voorbeeld wordt een teller per sleutel geïnitialiseerd vanuit een bestaande Delta-tabel:
Python (op regels gebaseerd)
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...
)
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...
Gebruik transformWithState in Lakeflow-pijplijnen
Gebruik de transformWithState operator in Lakeflow-pijplijnen om willekeurige stateful logica in uw streamingpijplijnen te implementeren met behulp van Python.
Voer hiervoor de volgende stappen uit:
- Definieer het uitvoerschema en de stateful-processorlogica voor uw willekeurige stateful-transformaties. Zie Voorbeeld van stateful toepassingen.
- Maak een Lakeflow-pijplijnstroom die de
transformWithStateoperator aanroept op een DataFrame. Zie Zelfstudie: Uw eerste pijplijn maken met behulp van de Lakeflow Pipelines Editor. - Voer uw pijplijn uit en valideer de resultaten op de doeltabel of de doellocatie.
Zie voor een voorbeeld dat transformWithState gebruikt om sensor heartbeats te bewaken, Voorbeeld: Gebruik transformWithState om sensor heartbeats te bewaken.