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.
Asynkron bearbetning med
Important
Asynkron bearbetning för Python:s radbaserade transformWithState API är i Beta. Se Azure Databricks förhandsversioner.
Asynkron bearbetning finns tillgänglig i Databricks Runtime 19 och uppåt.
Python transformWithState stödjer asynkron bearbetning byggd på asyncio. Genom att köra tillståndsoperationer och användarlogik samtidigt över gruppering av nycklar och batcha interprocesskommunikation har asynkron bearbetning högre genomströmning än synkron bearbetning med endast mindre kodändringar. Denna genomströmningsvinst kräver inga asynkrona tredjepartsbibliotek. Avancerade användare kan ytterligare optimera sina applikationer med asynkrona programmeringsmönster och asynkrona bibliotek.
För att använda asynkron bearbetning, implementera en AsyncStatefulProcessor istället för den synkrona StatefulProcessor. API:et AsyncStatefulProcessor speglar det synkrona StatefulProcessor API:et, så de flesta applikationer kräver endast små ändringar för att använda det asynkrona API:et. Se Implementera en AsyncStatefulProcessor.
För det synkrona transformWithState API:et och kärnkoncepten, se Bygg en anpassad tillståndsbaserad applikation med transformWithState.
Note
Asynkron bearbetning finns endast tillgänglig för Python:s radbaserade transformWithState API. Det stöds inte för transformWithStateInPandas eller för Scala transformWithState API. Asynkron bearbetning stöds inte i serverlös beräkning.
Implementera en AsyncStatefulProcessor
För att konvertera en synkron StatefulProcessor till en AsyncStatefulProcessor, gör följande ändringar:
- Definiera API-metoderna (
init, , ,handleExpiredTimerhandleInputRows, , ochhandleInitialState) med nyckelordetasync defclose. - Läs och uppdatera tillstånds- och timervärden med
await, eller kör dem med Python:sasynciobibliotek. Detta gäller för tillståndsoperationer såsomvalueState.get()och för timeroperationer såsomregisterTimer. Att skapa tillståndsobjekt, såsomhandle.getValueState, förblir synkront.
Följande överväganden gäller för asynkron bearbetning:
- Om din applikation lagrar data i medlemsvariabler eller i externa system rekommenderar Databricks att du skriver om logiken för att vara säker för samtidig exekvering. Eftersom
handleInputRowsochhandleExpiredTimerkan köras samtidigt över grupperingsnycklar får interleaved-körningar inte korrupta delad data. De flesta ansökningar uppfyller redan detta krav. - Databricks rekommenderar att du inte fångar upp eller undertrycker fel från tillståndsoperationer. Apache Spark hanterar dessa fel åt dig. Om en tillståndsoperation misslyckas, misslyckas uppgiften med Apache Spark och försöker igen.
- I ett
AsyncStatefulProcessortillstånd hanteras operationsfel åt dig och dyker aldrig upp i din kod. - I en synkron
StatefulProcessor, tillståndsdriftfel uppstår i din kod, men att undertrycka dem kan äventyra datakorrektheten.
- I ett
Exempel: räkna rader för varje grupperingsnyckel
Följande exempel definierar en AsyncCountProcessor som räknar antalet rader för varje grupperingsnyckel. Variabeln value_schema definierar schemat för som ValueState lagrar den löpande räkningen. Jämfört med en synkron StatefulProcessorär ändringarna nyckelordet async def för varje metod samt await för tillståndsläsning och uppdateringsoperationer. Samtalet till getValueState in init förblir synkront. Definiera processorn enligt följande kod:
from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor, AsyncStatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType
value_schema = StructType([StructField("count", LongType(), True)])
class AsyncCountProcessor(AsyncStatefulProcessor):
async def init(self, handle: AsyncStatefulProcessorHandle) -> None:
self.count = handle.getValueState("count", value_schema)
async def handleInputRows(self, key, rows, timerValues):
total = (await self.count.get() or (0,))[0]
for _ in rows:
total += 1
await self.count.update((total,))
yield Row(action=key[0], count=total)
async def close(self) -> None:
pass
Kör en fråga med en asynkron processor
För att köra en fråga med en asynkron processor, skicka din AsyncStatefulProcessor till transformWithState. Frågan använder samma syntax som den synkrona vägen. De asynkrona och synkrona API:erna delar samma tillståndsformat, så du kan växla en befintlig fråga mellan en AsyncStatefulProcessor och en synkron StatefulProcessor medan du återanvänder samma kontrollpunkt.
Exempel: räkna händelser i exempeldatasetet events
Följande exempel körs AsyncCountProcessorevents mot exempeldatasetet. Varje post har ett time fält (epoksekunder) och ett action fält med värdet Open eller Close. Frågegruppen grupperar efter action och räknar händelserna för varje åtgärdstyp. För fler exempeldataset, se Exempeldataset.
Variabeln input_schema definierar schemat för källpostarna, och variabeln output_schema definierar schemat för de rader som processorn genererar. För att läsa exempeldatasetet som en ström, definiera båda schemana och starta sedan frågan som i följande kod:
from pyspark.sql.types import StructType, StructField, StringType, LongType
input_schema = StructType([
StructField("time", LongType(), True),
StructField("action", StringType(), True),
])
output_schema = StructType([
StructField("action", StringType(), True),
StructField("count", LongType(), True),
])
events = (
spark.readStream.schema(input_schema)
.option("maxFilesPerTrigger", 10)
.json("/databricks-datasets/structured-streaming/events")
)
q = (
events.groupBy("action")
.transformWithState(
statefulProcessor=AsyncCountProcessor(),
outputStructType=output_schema,
outputMode="Update",
timeMode="None",
)
.writeStream.format("memory")
.queryName("async_counts")
.trigger(availableNow=True)
.start()
)
q.awaitTermination()
När frågan är klar, se den löpande räkningen för varje åtgärdstyp som i följande kod:
display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))
Asynkrona tillstånds- och timeroperationer
I en AsyncStatefulProcessorär tillståndsvariabel- och timeroperationer som läser eller skriver värden asynkrona. De flesta av dessa operationer returnerar ett enda resultat som du hämtar med await. Operationer som returnerar en samling returnerar istället en asynkron iterator som du konsumerar med async for. För en introduktion till async/await och asynkrona iteratorer i Python, se Python asyncio-dokumentationen.
Följande tabell listar operationer som returnerar ett enda resultat som du kan hämta med await:
| Class | Operationer som använder await |
|---|---|
AsyncValueState |
exists, get, update, clear |
AsyncMapState |
exists, getValue, containsKey, updateValue, , removeKey, clear |
AsyncListState |
exists, put, appendValue, appendList, clear |
AsyncStatefulProcessorHandle |
registerTimer, deleteTimer |
Följande tabell listar operationer som returnerar en asynkron iterator som du kan hämta med async for:
| Class | Operationer som använder async for |
|---|---|
AsyncMapState |
iterator, keys, values |
AsyncListState |
get |
AsyncStatefulProcessorHandle |
listTimers |
Exempel: async for
Till exempel, för att läsa värdena i en AsyncListState, iterera med async for som i följande kod:
total = 0
async for value in self.items.get():
total += value[0]
Metoderna som skapar tillståndsobjekt och tar bort tillståndsvariabler förblir synkrona: getValueState, getMapState, , getListStateoch deleteIfExists.
För en beskrivning av varje tillståndstyp, se Anpassade tillståndstyper.
Optimera med asynkrona programmeringsmönster
Asynkron bearbetning är användbar när din logik väntar på externa operationer, såsom nätverksförfrågningar. Istället för att vänta på varje förfrågan i följd, använd asyncio för att köra förfrågningarna samtidigt och minska vilotiden.
Exempel: kör samtidiga förfrågningar med asyncio.gather
Följande exempel använder asyncio.gather för att avfyra alla HTTP-förfrågningar per rad samtidigt och vänta på att de ska slutföras, och lagrar sedan maxpoängen i tillståndet. Definiera processorn enligt följande kod:
import asyncio
import aiohttp
from pyspark.sql import Row
from pyspark.sql.streaming import AsyncStatefulProcessor
class HttpScoreRowGatherProcessor(AsyncStatefulProcessor):
async def init(self, handle):
self._score_state = handle.getValueState("last_score", "score double")
self._session = aiohttp.ClientSession()
async def _fetch_score(self, row) -> float:
async with self._session.get(
f"https://api.example.com/score/{row.event_id}"
) as resp:
return (await resp.json())["score"]
async def handleInputRows(self, key, rows, timerValues):
user_id = key[0]
scores = await asyncio.gather(*[self._fetch_score(row) for row in rows])
max_score = max(scores)
await self._score_state.update((max_score,))
yield Row(user_id=user_id, score=max_score)
async def close(self):
await self._session.close()