Asynchrone verwerking met transformWithState (Beta)

Important

Asynchrone verwerking voor de Python-rijgebaseerde transformWithState API is in bèta. Zie Azure Databricks preview-releases.

Asynchrone verwerking is beschikbaar in Databricks Runtime 19 en hoger.

Python transformWithState ondersteunt asynchrone verwerking gebouwd op asyncio. Door toestandsoperaties en gebruikerslogica gelijktijdig te laten uitvoeren over het groeperen van sleutels en het batchen van interprocescommunicatie, heeft asynchrone verwerking een hogere doorvoer dan synchrone verwerking met slechts kleine codewijzigingen. Deze doorvoerwinst vereist geen asynchrone bibliotheken van derden. Gevorderde gebruikers kunnen hun applicaties verder optimaliseren met asynchroon programmeerpatronen en asynchroon geschikte bibliotheken.

Om asynchrone verwerking te gebruiken, implementeer je een AsyncStatefulProcessor in plaats van de synchrone StatefulProcessor. De AsyncStatefulProcessor API spiegelt de synchrone StatefulProcessor API, dus de meeste applicaties vereisen slechts kleine wijzigingen om de asynchrone API te gebruiken. Zie Implement an AsyncStatefulProcessor.

Voor de synchrone transformWithState API en kernconcepten, zie Bouw een aangepaste stateful applicatie met transformWithState.

Note

Asynchrone verwerking is alleen beschikbaar voor de Python-rijgebaseerde transformWithState API. Het wordt niet ondersteund voor transformWithStateInPandas of voor de Scala transformWithState API. Asynchrone verwerking wordt niet ondersteund in serverless compute.

Implementeer een AsyncStatefulProcessor

Om een synchrone StatefulProcessor om te zetten naar een AsyncStatefulProcessor, voer je de volgende wijzigingen aan:

  • Definieer de API-methoden (init, , closehandleInputRows, handleExpiredTimer, en handleInitialState) met het async def trefwoord.
  • Lees en werk status- en timerwaarden bij met await, of voer ze uit met de asyncio bibliotheek van Python. Dit geldt voor toestandsbewerkingen zoals valueState.get() en voor timerbewerkingen zoals registerTimer. Het creëren van toestandsobjecten, zoals handle.getValueState, blijft synchroon.

De volgende overwegingen gelden voor asynchrone verwerking:

  • Als je applicatie data opslaat in lidvariabelen of in externe systemen, raadt Databricks aan om de logica te herschrijven zodat deze veilig is voor gelijktijdige uitvoering. Omdat handleInputRows en handleExpiredTimer gelijktijdig kan draaien over groepsleutels, mogen interleaved runs gedeelde data niet beschadigen. De meeste sollicitaties voldoen al aan deze vereiste.
  • Databricks raadt aan om fouten uit toestandsoperaties niet te vangen of te onderdrukken. Apache Spark regelt deze fouten voor je. Als een toestandsoperatie faalt, faalt Apache Spark de taak en probeert het opnieuw.
    • In een AsyncStatefulProcessor, staat worden operatiefouten voor je beheerd en komen ze nooit naar voren in je code.
    • In een synchrone StatefulProcessortoestand worden bedrijfsfouten in je code geactiveerd, maar het onderdrukken ervan kan de correctheid van de data in gevaar brengen.

Voorbeeld: tel rijen voor elke groeperingssleutel

Het volgende voorbeeld definieert een AsyncCountProcessor die het aantal rijen voor elke groepssleutel telt. De value_schema variabele definieert het schema van de ValueState dat de lopende telling opslaat. In vergelijking met een synchrone StatefulProcessor, zijn de wijzigingen het async def sleutelwoord op elke methode en await op de toestandslees- en updateoperaties. De oproep naar getValueState binnen init blijft synchroon. Definieer de processor als in de volgende code:

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

Voer een query uit met een asynchroon processor

Om een query met een asynchroon processor uit te voeren, geef je je door AsyncStatefulProcessor aan transformWithState. De query gebruikt dezelfde syntaxis als het synchrone pad. De asynchrone en synchrone API's delen hetzelfde statusformaat, dus je kunt een bestaande query wisselen tussen een AsyncStatefulProcessor en een synchrone StatefulProcessor terwijl je hetzelfde checkpoint hergebruikt.

Voorbeeld: tel gebeurtenissen in de events voorbeelddataset

Het volgende voorbeeld draait AsyncCountProcessor tegen de events voorbeelddataset. Elke record heeft een time veld (epochseconden) en een action veld met de waarde Open of Close. De query groepeert op action en telt de gebeurtenissen voor elk actietype. Voor meer voorbeelddatasets, zie Voorbeelddatasets.

De input_schema variabele definieert het schema van de bronrecords, en de output_schema variabele definieert het schema van de rijen die de processor uitzendt. Om de voorbeelddataset als een stroom te lezen, definieer je beide schema's en start je de query zoals in de volgende code:

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

Na voltooiing van de query bekijk je de lopende telling voor elk actietype zoals in de volgende code:

display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))

Asynchrone toestands- en timeroperaties

In een AsyncStatefulProcessor, zijn toestandsvariabelen en timeroperaties die waarden lezen of schrijven asynchroon. De meeste van deze bewerkingen geven één resultaat dat je ophaalt met await. Bewerkingen die een collectie teruggeven, geven in plaats daarvan een asynchrone iterator die je consumeert met async for. Voor een introductie tot async/await en asynchrone iterators in Python, zie de Python asyncio-documentatie.

De volgende tabel geeft bewerkingen weer die één enkel resultaat teruggeven waarmee awaitu kunt ophalen:

Class Bewerkingen die gebruik maken van await
AsyncValueState exists,get,update,clear
AsyncMapState exists, getValue, containsKey, updateValue, , removeKey, clear
AsyncListState exists, put, appendValue, appendList, , clear
AsyncStatefulProcessorHandle registerTimer, deleteTimer

De volgende tabel geeft bewerkingen weer die een asynchrone iterator teruggeven waarmee async forje kunt ophalen:

Class Bewerkingen die gebruik maken van async for
AsyncMapState iterator, keys, values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

Voorbeeld:async for

Om bijvoorbeeld de waarden in een AsyncListStatete lezen, itereren met async for zoals in de volgende code:

total = 0
async for value in self.items.get():
  total += value[0]

De methoden die toestandsobjecten creëren en toestandsvariabelen verwijderen blijven synchroon: getValueState, getMapState, , getListStateen deleteIfExists.

Voor een beschrijving van elk toestandstype, zie Aangepaste toestandtypes.

Optimaliseer met asynchroon programmeerpatronen

Asynchrone verwerking is nuttig wanneer je logica wacht op externe bewerkingen, zoals netwerkverzoeken. In plaats van op elk verzoek in volgorde te wachten, gebruik asyncio je om de verzoeken gelijktijdig uit te voeren en de inactieve tijd te verminderen.

Voorbeeld: voer gelijktijdige verzoeken uit met asyncio.gather

Het volgende voorbeeld gebruikt asyncio.gather om alle HTTP-verzoeken per rij gelijktijdig af te vuren en te wachten tot ze voltooid zijn, waarna de maximale score in de toestand wordt opgeslagen. Definieer de processor als in de volgende code:

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