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.
Asynchrone verwerking met
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, enhandleInitialState) met hetasync deftrefwoord. - Lees en werk status- en timerwaarden bij met
await, of voer ze uit met deasynciobibliotheek van Python. Dit geldt voor toestandsbewerkingen zoalsvalueState.get()en voor timerbewerkingen zoalsregisterTimer. Het creëren van toestandsobjecten, zoalshandle.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
handleInputRowsenhandleExpiredTimergelijktijdig 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.
- In een
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()