Asynkron bearbetning med transformWithState (Beta)

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, , och handleInitialState) med nyckelordet async defclose.
  • Läs och uppdatera tillstånds- och timervärden med await, eller kör dem med Python:s asyncio bibliotek. Detta gäller för tillståndsoperationer såsom valueState.get() och för timeroperationer såsom registerTimer. Att skapa tillståndsobjekt, såsom handle.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 handleInputRows och handleExpiredTimer kan 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.

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()