Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Bygg en anpassad tillståndsbaserad applikation med
Du kan använda transformWithState för att skapa tillståndskänsliga strömningsprogram och för att implementera lösningar med låg svarstid och nära realtid. Med anpassade tillståndskänsliga operatorer kan du skapa godtycklig tillståndskänslig logik som gör att du kan skapa nya användningsfall som inte är möjliga med traditionell bearbetning av strukturerad direktuppspelning.
Anteckning
För tillståndskänsliga åtgärder som aggregeringar, deduplicering och strömningsanslutningar rekommenderar Databricks att du använder inbyggda operatorer för strukturerad direktuppspelning i stället för anpassad logik. Se Vad är tillståndskänslig strömning?.
Databricks rekommenderar att du använder transformWithState i stället för äldre operatorer, till exempel flatMapGroupsWithState och mapGroupsWithState, för godtyckliga tillståndstransformeringar. Se Äldre godtyckliga tillståndskänsliga operatorer.
Krav
Operatorerna transformWithState och transformWithStateInPandas har följande krav:
- Tillgänglig i Databricks Runtime 16.2 och senare.
- I realtidsläge använder du Databricks Runtime 17.3 LTS eller senare. Se koncept för realtidsläge.
- För standardåtkomstläge är Python tillgängligt i Databricks Runtime 16.3 och senare, och Scala är tillgängligt i Databricks Runtime 17.3 och senare.
- RocksDB är standardtillståndslagringsprovidern i Databricks Runtime 17.3 och senare.
För Databricks Runtime 17.2 och nedan måste du konfigurera RocksDB-tillståndslagerprovidern. Databricks rekommenderar att du aktiverar RocksDB i Spark-konfigurationen.
spark.conf.set("spark.sql.streaming.stateStore.providerClass", "org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider")
Vad är transformWithState?
Operatorn transformWithState tillämpar en anpassad tillståndskänslig processor på en fråga för strukturerad direktuppspelning. Du måste implementera en anpassad tillståndskänslig processor för att kunna använda transformWithState. Strukturerad direktuppspelning innehåller API:er för att skapa en tillståndskänslig processor med hjälp av Python, Scala eller Java.
Använd transformWithState för att tillämpa anpassad logik på en grupperingsnyckel. Följande beskriver designen på hög nivå:
- Definiera en eller flera tillståndsvariabler.
- Tillståndsinformationen bevaras för varje grupperingsnyckel. Du kan komma åt varje tillståndsvariabel i användardefinierad kod.
- För varje mikrobatch som bearbetas finns alla rader för nyckeln tillgängliga i form av en iterator.
- Använd
StatefulProcessorHandletillsammans med timrar och användardefinierade villkor för att styra hur rader ska matas ut. - För att hantera tillståndsförfallodatum och tillståndsstorlek har tillståndsvärden stöd för enskilda TTL-definitioner (time to live).
Eftersom transformWithState stöder schemautveckling i tillståndsarkivet kan du iterera och uppdatera dina produktionsprogram utan att förlora information om historiska tillstånd. När du har uppdaterat tillståndsschemat behöver du inte bearbeta rader igen, vilket förenklar koddistributioner och underhåll. Se även Schemautveckling i tillståndslagret.
Viktig
Azure Databricks dokumentation används transformWithState för att beskriva både Python- och Scala-implementeringar:
- PySpark stöder både det radbaserade
transformWithStateAPI:et och den Pandas-baseradetransformWithStateInPandasoperatorn.-
transformWithStateInPandasstöds inte i realtidsläge. AnvändtransformWithStatei stället . Mer information finnstransformWithStatei realtidsläge. - Det radbaserade
transformWithStateAPI:et stöder asynkron bearbetning medasyncioför högre genomströmning. Asynkron bearbetning stöds inte i serverlös beräkning. Se Asynkron bearbetning (Beta).
-
- Scala stöder endast det radbaserade
transformWithStateAPI:et.
Scala- och Python implementeringar av transformWithState har samma funktioner, men med vissa skillnader i syntax.
Att definiera en StatefulProcessor
Du definierar en tillståndskänslig processor genom att StatefulProcessor utöka klassen och implementera dess metoder.
Spark skickar en StatefulProcessorHandle till metoden init i din StatefulProcessor. Använd handtaget för att skapa tillståndsvariabler och interagera med tillståndsarkivet.
transformWithState stöder tre tillståndstyper: ValueState, ListStateoch MapState. Varje typ lagrar tillstånd för varje grupperingsnyckel med hjälp av en annan underliggande datastruktur.
Implementera följande metoder för att definiera din anpassade logik:
- Implementera
handleInputRowsför att styra hur programmet bearbetar data, uppdaterar tillstånd och genererar rader för varje mikrobatch. Se att hantera indatarader. - Implementera
handleExpiredTimerför att köra tidsbaserad logik oavsett om grupperingsnyckeln tar emot nya rader i en mikrobatch. Se Hantera utlöpta timerar. - Du kan också implementera
handleInitialStateför att fylla i tillståndet i förväg innan programmet bearbetar indatarader. Se Hantera inledande tillstånd.
I följande tabell jämförs de funktionella beteendena för dessa metoder:
| Uppförande | handleInputRows |
handleExpiredTimer |
|---|---|---|
| Hämta, placera, uppdatera eller rensa tillståndsvärden | Ja | Ja |
| Skapa eller ta bort en timer | Ja | Ja |
| Mata ut rader | Ja | Ja |
| Iterera över rader i den aktuella mikrobatchen | Ja | Nej |
| Utlösarlogik baserat på förfluten tid | Nej | Ja |
Du kan kombinera båda handleInputRows och handleExpiredTimer implementera komplex logik efter behov.
Du kan till exempel implementera ett program som använder handleInputRows för att uppdatera tillståndsvärden för varje mikrobatch och ange en timer på 10 sekunder i framtiden. Om inga ytterligare rader bearbetas, kan du använda handleExpiredTimer för att skicka ut de aktuella värdena i tillståndslagret. Om nya rader bearbetas för grupperingsnyckeln kan du rensa den befintliga timern och ange en ny timer.
StatefulProcessorHandle
I PySpark StatefulProcessorHandle låter klassen dig komma åt funktioner som styr hur koden använder tillståndsinformation.
När du initierar en StatefulProcessor måste du alltid importera StatefulProcessorHandle och skicka den till variabeln handle. Variabeln handle kopplar den lokala variabeln i klassen Python till tillståndsvariabeln.
Anteckning
Scala använder metoden getHandle.
Anpassade tillståndstyper
Du kan implementera flera tillståndsobjekt i en enda tillståndskänslig operator.
Välj en tillståndstyp baserat på din fullständiga programlogik. Du kan till exempel spåra sessioner med en ValueState grupperad efter user_id och session_id. Eller, för att utvärdera villkor över flera sessioner, använder du en MapState grupperad efter user_id med session_id som mappningsnyckel.
Om tillståndsobjektet använder ett StructTypemåste du definiera unika namn för varje fält i schemats struktur. Dessa namn visas när du läser tillståndsarkivet. Se läsa information om status för strukturerad direktuppspelning.
I följande avsnitt beskrivs de tillståndstyper som stöds av transformWithState:
ValueState
ValueState lagrar ett värde för varje grupperingsnyckel.
Ett värde kan innehålla komplexa typer, till exempel en struct eller en tupl. För ValueStatemåste du implementera logik för att ersätta hela värdet.
Livslängden för ett värdetillstånd återställs när värdet uppdateras. Om du bearbetar en källnyckel för ValueState utan att uppdatera den lagrade ValueStateåterställs inte time-to-live.
ListState
ListState lagrar en lista för varje grupperingsnyckel.
Ett listtillstånd är en samling värden som var och en kan innehålla komplexa typer. Varje värde i en lista har sin egen time-to-live.
Du kan lägga till objekt i en lista genom att lägga till enskilda objekt, lägga till en lista med objekt eller skriva över hela listan med en put. Om du vill återställa time-to-live måste du använda en put åtgärd.
MapState
MapState lagrar en karta för varje grupperingsnyckel. Kartor är Apache Spark-motsvarigheten till en Python ordlista (dict).
Ett mappningstillstånd är en samling distinkta nycklar där varje nyckel motsvarar ett värde, och varje värde kan innehålla komplexa typer. Varje nyckel/värde-par i en mappning har en egen livslängd.
Du kan uppdatera värdet för en specifik nyckel eller ta bort en nyckel och dess värde. Du kan returnera ett enskilt värde med hjälp av dess nyckel, visa alla nycklar, visa alla värden eller returnera en iterator för att arbeta med den fullständiga uppsättningen nyckel/värde-par på kartan.
Viktig
Grupperingsnycklar beskriver de fält som anges i GROUP BY-satsen i frågan Strukturerad direktuppspelning. Map-tillstånd kan innehålla ett godtyckligt antal nyckel/värde-par för en grupperingsnyckel.
Om din fråga till exempel använder GROUP BY user_id och du vill definiera en karta för varje session_id, är user_id din grupperingsnyckel och MapState nyckeln är session_id:
Python
class SessionTracker(StatefulProcessor):
def init(self, handle: StatefulProcessorHandle) -> None:
self.sessions = handle.getMapState("sessions", "session_id string", "count long")
def handleInputRows(self, key, rows: Iterator[Row], timerValues) -> Iterator[Row]:
for row in rows:
session_key = (row["session_id"],) # session_id is the MapState key
count = self.sessions.getValue(session_key)[0] if self.sessions.containsKey(session_key) else 0
new_count = count + 1
self.sessions.updateValue(session_key, (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, (String, Long)] {
@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[(String, Long)] = {
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())
Skapa en anpassad tillståndsvariabel i StatefulProcessor
När du initierar din StatefulProcessorskapar du en lokal variabel för varje tillståndsobjekt som gör att du kan interagera med tillståndsobjekt i din anpassade logik. Definiera och initiera tillståndsvariabler genom att åsidosätta den inbyggda init metoden i StatefulProcessor klassen.
Du kan definiera valfritt antal tillståndsobjekt med metoderna getValueState, getListStateoch getMapState i .StatefulProcessor
Varje tillståndsobjekt måste ha följande:
- Ett unikt namn
- Ett schema
- I Python måste du ange schemat.
- I Scala kan du skicka ett
Encoderför att ange tillståndsschema.
Du kan också ange en TTL-varaktighet (time-to-live) i millisekunder. Om du implementerar ett karttillstånd måste du ange en separat schemadefinition för kartnycklarna och värdena.
Anteckning
StatefulProcessor hanterar logik separat för frågor, uppdateringar och utsändning av tillståndsinformation. Se Använda dina tillståndsvariabler i metoder med anpassad logik.
Använda dina tillståndsvariabler i metoder med anpassad logik
Tillståndsobjekt har metoder för att hämta tillstånd, uppdatera befintlig tillståndsinformation och rensa det aktuella tillståndet.
Varje grupperingsnyckel har information om dedikerat tillstånd.
-
StatefulProcessorgenererar rader baserat på din anpassade logik och det angivna utdataschemat. Se Skapa rader. - Använd
statestore-läsaren för att komma åt värden i tillståndslagret. Den här läsaren är avsedd för batcharbetsbelastningar och är inte avsedd för arbetsbelastningar med låg fördröjning. Se läsa information om status för strukturerad direktuppspelning. - Logik som endast anges med
handleInputRowskörs bara om det finns rader för nyckeln i en mikrobatch. Se att hantera indatarader. - Använd
handleExpiredTimerför att implementera tidsbaserad logik som inte är beroende av att övervaka rader för att utlösas. Se Hantera utlöpta timerar.
Anteckning
Tillståndsobjekt isoleras genom gruppering av nycklar med följande konsekvenser:
- Tillståndsvärden kan inte påverkas av rader som är associerade med en annan grupperingsnyckel.
- Du kan inte implementera logik som är beroende av att jämföra värden eller uppdatera tillstånd mellan grupperingsnycklar.
Du kan jämföra värden i en grupperingsnyckel. Använd en MapState för att implementera logik med en andra nyckel som din anpassade logik kan använda. Om du till exempel grupperar user_id efter och använder ip_address för din MapState nyckel kan du spåra samtidiga användarsessioner.
Avancerade överväganden för att arbeta med state
Tillståndsuppdateringar är feltoleranta. Om en aktivitet kraschar innan en mikrobatch har slutfört bearbetningen använder återförsöket värdet från den senaste lyckade mikrobatchen.
För optimerad prestanda rekommenderar Databricks att du bearbetar alla värden i iteratorn för en viss nyckel och genomför uppdateringar i en enda skrivning. När du skriver till en tillståndsvariabel utlöser detta en skrivning till RocksDB.
Tillståndsvärden har inte standardvärden. Om logiken kräver att du läser befintlig tillståndsinformation använder du exists metoden.
Om du vill implementera logik för null-tillstånd MapState kan du med variabler söka efter enskilda nycklar eller visa en lista över alla nycklar.
Hantera indatarader
handleInputRows Använd metoden för att definiera hur programmet bearbetar rader och uppdaterar tillståndsvärden. Den här metoden anropas varje gång din Structured Streaming-fråga bearbetar rader för en viss grupperingsnyckel.
För de flesta tillståndskänsliga program som implementeras med transformWithStatedefinieras kärnlogik med hjälp av handleInputRows.
För varje mikrobatchuppdatering som bearbetas är alla rader i mikrobatchen för en viss grupperingsnyckel tillgängliga med hjälp av en iterator. Användardefinierad logik kan interagera med alla rader från den aktuella mikrobatchen och värdena i tillståndsarkivet.
Hantera utgångna timrar
handleExpiredTimer Använd metoden för att implementera anpassad logik baserat på förfluten tid.
Inom en grupperingsnyckel identifieras timers unikt av tidsstämpeln.
När en timer upphör att gälla bestäms resultatet av logiken som implementeras i ditt program. Vanliga mönster är:
- Genererar information som lagras i en tillståndsvariabel.
- Rensa lagrad tillståndsinformation.
- Skapa en ny timer.
Utgångna timers utlöses även om inga rader för deras tillhörande nyckel bearbetas i en mikrobatch.
Ange tidsläget
När du skickar StatefulProcessor till transformWithStatemåste du ange tidsläget med hjälp av parametern timeMode .
Följande alternativ stöds:
| Tidsläge | Beskrivning |
|---|---|
ProcessingTime |
Både timrar och TTL stöds och utvärderas baserat på klocktiden vid bearbetningen av varje mikrobatch i Apache Spark. Använd ProcessingTime när du vill att timers ska utlösas med ett fast intervall i förhållande till när rader bearbetas, oavsett tidsstämplar i data. |
EventTime |
Timers stöds och utvärderas utifrån händelsetidsvattenstämpeln. Vattenstämpeln utvecklas när Apache Spark observerar tidsstämplar i indata. TTL stöds inte med EventTime. Använd EventTime när dina data innehåller tidsstämplar och du vill att timers ska utlösas baserat på förloppet för dessa tidsstämplar. När du använder EventTimemåste du också ange parametern eventTimeColumnName . Se även eventTimeColumnName. |
NoTime eller TimeMode.None() |
Timers och TTL stöds inte. Använd NoTime när ditt tillståndskänsliga program inte kräver tidsbaserad logik. |
eventTimeColumnName
När du använder EventTime tidsläget anger parametern eventTimeColumnName namnet på kolumnen i utdataschemat som innehåller händelsetidsstämpeln. Apache Spark använder den här kolumnen för att vidarebefordra vattenstämpeln till utdataströmmen, vilket möjliggör korrekta efterföljande tidsbaserade operationer.
Python
eventTimeColumnName är ett ytterligare argument för transformWithState eller transformWithStateInPandas:
q = (
df.groupBy("key")
.transformWithState(
statefulProcessor=MyProcessor(),
outputStructType=output_schema,
outputMode="Append",
timeMode="EventTime",
eventTimeColumnName="outputTimestamp",
)
.writeStream...
)
Scala
transformWithState accepterar eventTimeColumnName i stället för timeMode. Den här metoden använder alltid läget EventTime:
val q = spark
.readStream
.format("delta")
.load(srcDeltaTableDir)
.as[(String, String)]
.groupByKey(x => x._1)
.transformWithState(
new MyProcessor(),
"outputTimestamp",
OutputMode.Append(),
)
.writeStream...
Inbyggda timervärden
Databricks rekommenderar att du inte anropar systemklockan i ditt anpassade tillståndskänsliga program, eftersom detta kan leda till otillförlitliga återförsök vid aktivitetsfel. Använd metoderna i klassen TimerValues när du måste komma åt bearbetningstiden eller vattenstämpeln:
TimerValues |
Beskrivning |
|---|---|
getCurrentProcessingTimeInMs |
Returnerar tidsstämpeln för bearbetningstiden för den aktuella batchen i millisekunder sedan epoken. |
getCurrentWatermarkInMs |
Returnerar tidsstämpeln för den aktuella batchens vattenstämpel i millisekunder sedan epok. |
Anteckning
Bearbetningstiden beskriver den tid då mikrobatchen bearbetas av Apache Spark. Många strömmande källor, till exempel Kafka, inkluderar även systembearbetningstid.
Vattenstämplar i strömmande sökfrågor definieras ofta mot evenemangstid eller bearbetningstid för den strömmande källan. Se Använd vattenstämplar för att kontrollera tröskelvärden för databehandling.
Både vattenstämplar och fönster kan användas i kombination med transformWithState. Du kan implementera liknande funktioner i ditt anpassade tillståndskänsliga program genom att använda TTL, timers och MapState eller ListState funktioner.
Time-to-live (TTL) för tillståndstyper
För att förhindra out-of-memory-fel och för att ta bort värden för inaktuell tillståndstyp har transformWithState stöd för ett valfritt TTL-värde (time to live) för varje tillståndstypvärde. När giltighetstiden har löpt ut tar TTL bort värden av tillståndstyp utan avisering. TTL kör inte handleExpiredTimer eller någon anpassad logik. Om du vill köra kod när tillståndet upphör att gälla använder du en timer i stället.
Viktig
Om du inte inför TTL måste du hantera utrensning av tillstånd för att undvika fel på grund av minnesbrist.
För alla tillståndstyper återställs TTL vid uppdatering av tillståndsinformation. TTL tillämpas för varje tillståndstypvärde, med olika regler för varje tillståndstyp:
- Tillståndsvariabler är begränsade till grupperingsnycklar.
- För
ValueStateobjekt lagras endast ett enda värde per grupperingsnyckel. TTL gäller för det här värdet. - För
ListStateobjekt kan listan innehålla många värden. TTL gäller för varje värde i en lista oberoende av varandra.- Även om TTL är begränsat till enskilda värden i en
ListState, är det enda sättet att uppdatera ett enskilt värde medputmetoden, som skriver över hela innehållet i variabelnListStateoch återställer TTL för alla värden i listan.
- Även om TTL är begränsat till enskilda värden i en
- För
MapStateobjekt har varje kartnyckel ett associerat tillståndsvärde. TTL gäller oberoende av varje nyckel/värde-par på en karta.
Anteckning
Med timers kan du definiera anpassad logik bortom tillståndsavhysning, inklusive att generera rader. Du kan också använda timers för att både rensa tillståndsinformation för ett visst tillståndsvärde och generera värden eller utlösa villkorslogik. Se Hantera utlöpta timerar.
Exempel på tillståndskänsligt program
I följande exempel definieras en anpassad tillståndskänslig processor, SimpleCounterProcessor, inklusive exempeltillståndsvariabler.
SimpleCounterProcessor använder ValueState, ListStateoch MapState för att räkna rader för varje grupperingsnyckel.
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]:
# Seed the running total from state so the count accumulates across micro-batches
count = self.value_state.get()[0] if self.value_state.exists() else 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 (radbaserad)
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]:
# Seed the running total from state so the count accumulates across micro-batches
count = self.value_state.get()[0] if self.value_state.exists() else 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[0], 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...
Kör exemplet från början till slut
Anteckning
De körbara exemplen på denna sida skapar tabeller i ett dedikerat main.stateful_examples schema så att de kan köras utan att påverka din befintliga data. Om du inte har behörighet att skapa scheman i katalogen main , ändra katalogen och schemat i exemplen till en plats där du kan skapa tabeller.
Processorn ovan definierar den tillståndsfulla logiken men startar inte en fråga. För att köra SimpleCounterProcessor med copy-paste, initiera en liten Delta Lake-tabell som källa för dataströmmen och starta sedan en fråga som skriver till en sink i minnet. Detta exempel använder Trigger.AvailableNow så att frågan bearbetar de seedade raderna och stannar. För att seeda källkoden och starta frågan, kör följande:
import uuid
# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")
# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_counter_source")
spark.createDataFrame(
[("a", "1"), ("a", "2"), ("a", "3"), ("b", "1"), ("b", "2")],
"key string, value string",
).write.saveAsTable("main.stateful_examples.tws_counter_source")
df = spark.readStream.table("main.stateful_examples.tws_counter_source")
q = (
df.groupBy("key")
.transformWithState(
statefulProcessor=SimpleCounterProcessor(),
outputStructType=output_schema,
outputMode="Update",
timeMode="None",
)
.writeStream.format("memory")
.queryName("counter_output")
.option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
.trigger(availableNow=True)
.start()
)
q.awaitTermination()
När frågan är klar, se räkningen för varje grupperingsnyckel:
display(spark.sql("SELECT id, countAsString FROM counter_output ORDER BY id"))
Nyckeln a har tre rader och nyckeln b har två, så frågan returnerar:
id countAsString
a 3
b 2
Fler exempel finns i Exempel på tillståndskänsliga program.
Anteckning
I Python är tillståndsvärden tupler. Skicka tupler till put och update, och förvänta dig tupler från get.
Om till exempel schemat för ditt ValueState är ett heltal:
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
Använd den här metoden för objekt i en ListState eller värden i en MapState också.
Mata ut rader
Du måste använda handleInputRows eller handleExpiredTimer för att definiera hur transformWithState genererar rader för varje grupperingsnyckel. Se Hantera indatarader och Hantera utgångna timers.
Anpassade tillståndsbaserade applikationer gör inga antaganden om hur tillståndsinformation används. För ett visst villkor kan programmet inte generera några rader, en rad eller många rader.
Anteckning
Du kan implementera flera tillståndsvärden och definiera flera villkor för att generera rader, men alla rader måste använda samma schema.
Python (Pandas)
Med transformWithStateInPandasdefinierar du utdataschemat med nyckelordet outputStructType .
Generera rader med hjälp av ett Pandas DataFrame-objekt och yield.
Du kan också använda yield en tom DataFrame. Om du använder update utdataläget och genererar en tom DataFrame uppdateras värdena för att grupperingsnyckeln ska vara null.
Python (radbaserad)
Med transformWithStatedefinierar du utdataschemat med nyckelordet outputStructType .
Generera rader med hjälp av ett Row objekt och yield.
Du kan också returnera en tom iterator. Om du använder update utdataläget och genererar en tom iterator uppdateras värdena för att grupperingsnyckeln ska vara null.
Scala
I Scala genererar du rader med hjälp av ett Iterator objekt. Schemat härleds automatiskt från schemat för de utgivna raderna.
Om du vill kan du returnera en tom Iterator. Om du använder update utdataläget och genererar ett tomt Iteratoruppdaterar detta värdena för att grupperingsnyckeln ska vara null.
Hantera inledande tillstånd
Du kan också skicka ett initialt tillstånd till den första mikrobatchen.
Du kan till exempel använda detta för att:
- Migrera ett befintligt arbetsflöde till ett nytt anpassat program.
- Uppgradera en tillståndskänslig operator för att ändra schemat eller logiken.
- Reparera ett fel som inte kan repareras automatiskt och som kräver manuella åtgärder.
Anteckning
Använd tillståndsarkivläsaren för att fråga efter tillståndsinformation från en befintlig kontrollpunkt. Se läsa information om status för strukturerad direktuppspelning.
Om du konverterar en befintlig Delta-tabell till ett tillståndskänsligt program läser du tabellen med hjälp av spark.read.table("table_name") och skickar den resulterande DataFrame.If you are converting an existing Delta table to a stateful application, read the table using spark.read.table("table_name") and pass the resulting DataFrame. Du kan också välja eller ändra fält så att de överensstämmer med ditt nya tillståndskänsliga program.
Du anger ett initialt tillstånd med hjälp av en DataFrame med samma grupperingsnyckelschema som indataraderna.
Anteckning
Python använder handleInitialState för att ange det ursprungliga tillståndet när du definierar ett StatefulProcessor. Scala använder den distinkta klassen StatefulProcessorWithInitialState.
I följande exempel används en räknare per nyckel från en befintlig Delta-tabell:
Python (radbaserad)
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),
])
import uuid
# Create a dedicated schema for the example tables
spark.sql("CREATE SCHEMA IF NOT EXISTS main.stateful_examples")
# Seed existing per-key counts to load as the initial state
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.existing_counts")
spark.createDataFrame(
[("x", 10)],
"id string, count int",
).write.saveAsTable("main.stateful_examples.existing_counts")
# Seed a small Delta table to use as the streaming source
spark.sql("DROP TABLE IF EXISTS main.stateful_examples.tws_initial_source")
spark.createDataFrame(
[("x", "a"), ("x", "b")],
"id string, value string",
).write.saveAsTable("main.stateful_examples.tws_initial_source")
df = spark.readStream.table("main.stateful_examples.tws_initial_source")
# Load existing counts as initial state — must use the same grouping key as the input
initial_state = spark.read.table("main.stateful_examples.existing_counts").groupBy("id")
q = (
df.groupBy("id")
.transformWithState(
statefulProcessor=CounterWithInitialState(),
outputStructType=output_schema,
outputMode="Update",
timeMode="None",
initialState=initial_state,
)
.writeStream.format("memory")
.queryName("initial_state_output")
.option("checkpointLocation", f"/tmp/checkpoint_{uuid.uuid4()}")
.trigger(availableNow=True)
.start()
)
q.awaitTermination()
# The initial state seeds "x" with 10, and the source adds two rows, so the count is 12
display(spark.sql("SELECT id, count FROM initial_state_output ORDER BY id"))
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...
Asynkron behandling (Beta)
Python transformWithState stödjer asynkron bearbetning med asyncio för att köra tillståndsoperationer och användarlogik samtidigt. Asynkron bearbetning har högre genomströmning än synkron bearbetning och kräver endast mindre kodändringar, utan några asynkrona tredjepartsbibliotek. För att använda asynkron bearbetning, implementera en AsyncStatefulProcessor istället för den synkrona StatefulProcessor. Se Asynkron bearbetning med transformWithState (Beta).
Använd transformWithState i Lakeflow-pipelines
Använd operatorn transformWithState i Lakeflow-pipelines för att implementera godtycklig tillståndskänslig logik i dina strömningspipelines med hjälp av Python.
Utför följande steg för att göra det:
- Definiera utdataschemat och tillståndskänslig processorlogik för dina godtyckliga tillståndskänsliga transformeringar. Exempel finns i Exempel på tillståndskänsliga program.
- Skapa ett Lakeflow-pipelineflöde som anropar operatorn
transformWithStatepå en DataFrame. Se Självstudie: Skapa din första pipeline med Lakeflow Pipelines-redigeraren. - Kör pipelinen och verifiera resultatet på måltabellen eller data-sänkan.
Ett exempel som använder transformWithState för att övervaka sensor pulsslag finns i Exempel: Använd transformWithState för att övervaka sensor pulsslag.