Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Przetwarzanie asynchroniczne z
Important
Asynchroniczne przetwarzanie dla API opartego transformWithState na wierszach Python jest w fazie beta. Zobacz Azure Databricks wersje zapoznawcze.
Przetwarzanie asynchroniczne jest dostępne w Databricks Runtime 19 i wyższych.
Python transformWithState obsługuje asynchroniczne przetwarzanie oparte na asyncio. Poprzez jednoczesne uruchamianie operacji stanu i logiki użytkownika przez grupowanie kluczy oraz grupowanie komunikacji międzyprocesowej, przetwarzanie asynchroniczne osiąga wyższą przepustowość niż synchroniczne, przy jedynie drobnych zmianach w kodzie. To wzmocnienie przepustowości nie wymaga żadnych zewnętrznych bibliotek asynchronicznych. Zaawansowani użytkownicy mogą dodatkowo optymalizować swoje aplikacje dzięki wzorcom programowania asynchronicznego i bibliotekom wspierającym asynchronię.
Aby użyć przetwarzania asynchronicznego, zaimplementuj an AsyncStatefulProcessor zamiast synchronicznego StatefulProcessor. API odzwierciedla AsyncStatefulProcessor synchroniczne StatefulProcessor API, więc większość aplikacji wymaga jedynie drobnych zmian, aby korzystać z asynchronicznego API. Zobacz Implementuj .AsyncStatefulProcessor
Dla synchronous transformWithState API i podstawowych koncepcji, zobacz Zbuduj niestandardową aplikację stanową z .transformWithState
Uwaga / Notatka
Przetwarzanie asynchroniczne jest dostępne tylko dla API opartego transformWithState na wierszach Python. Nie jest obsługiwany ani transformWithStateInPandas wspierany dla API Scala transformWithState . Przetwarzanie asynchroniczne nie jest obsługiwane w przetwarzaniu serwerowym.
Implementuj AsyncStatefulProcessor
Aby przekształcić syntrybuł StatefulProcessor w , AsyncStatefulProcessorwykonaj następujące zmiany:
- Zdefiniuj metody API (
init,close,handleInputRows,handleExpiredTimer, orazhandleInitialState) za pomocąasync defsłowa kluczowego. - Czytaj i aktualizuj wartości stanu i timera za pomocą
await, albo uruchamiajasyncioje za pomocą biblioteki Python. Dotyczy to operacji stanów, takich jak orazvalueState.get()operacji timera, takich jakregisterTimer. Tworzenie obiektów stanów, takich jakhandle.getValueState, pozostaje synchroniczne.
Następujące kwestie dotyczą przetwarzania asynchronicznego:
- Jeśli Twoja aplikacja przechowuje dane w zmiennych członkowskich lub w zewnętrznych systemach, Databricks zaleca przepisanie logiki, aby była bezpieczna do współbieżnego wykonywania. Ponieważ
handleInputRowsi mogąhandleExpiredTimerdziałać równocześnie na kluczach grupowania, przeplatane przebiegi nie mogą uszkadzać danych współdzielonych. Większość aplikacji już spełnia ten wymóg. - Databricks zaleca, aby nie wykrywać ani nie tłumić błędów w operacjach stanowych. Apache Spark radzi sobie z tymi błędami za Ciebie. Jeśli operacja stanu nie powiodła się, Apache Spark nie wykonuje zadania i podejmuje próbę ponownie.
- W
AsyncStatefulProcessorstanie operacji , błędy są zarządzane za Ciebie i nigdy nie pojawiają się w twoim kodzie. - W synchronicznym
StatefulProcessorstanie , błędy operacji pojawiają się w kodzie, ale ich tłumienie może zagrozić poprawności danych.
- W
Przykład: liczenie wierszy dla każdego klucza grupowania
Poniższy przykład definiuje , AsyncCountProcessor która liczy liczbę wierszy dla każdego klucza grupowania. Zmienna definiuje value_schema schemat , ValueState który przechowuje liczbę bieżącą. W porównaniu do synchronicznego StatefulProcessor, zmiany to async def słowa kluczowe dla każdej metody oraz await w operacjach odczytu i aktualizacji stanu. Wezwanie do getValueState "wejść" init pozostaje synchroniczne. Zdefiniuj procesor w następującym kodzie:
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
Wykonaj zapytanie z procesorem asynchronicznym
Aby wykonać zapytanie z procesorem asynchronicznym, przekaż swoje AsyncStatefulProcessor do transformWithState. Zapytanie korzysta z tej samej składni co ścieżka synchroniczna. API asynchroniczne i synchroniczne mają ten sam format stanu, więc możesz przełączać istniejące zapytanie między an AsyncStatefulProcessor a synchroniczne StatefulProcessor , używając tego samego punktu kontrolnego.
Przykład: zliczenie zdarzeń w próbnym zbiorze events danych
Poniższy przykład jest połączany AsyncCountProcessor z przykładowym zbiorem events danych. Każdy rekord ma pole time (sekundy epoki) oraz action pole o wartości Open lub Close. Zapytanie grupuje według action i liczy zdarzenia dla każdego typu akcji. Więcej przykładowych zbiorów danych można znaleźć w sekcji Przykładowe zbiory danych.
Zmienna input_schema definiuje schemat rekordów źródłowych, a zmienna output_schema schemat wierszy generowanych przez procesor. Aby odczytać przykładowy zbiór danych jako strumień, zdefiniuj oba schematy, a następnie rozpocznij zapytanie zgodnie z następującym kodem:
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()
Po zakończeniu zapytania zobacz liczbę bieżących dla każdego typu akcji w następującym kodzie:
display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))
Operacje stanu asynchronicznego i timera
W AsyncStatefulProcessor, operacje zmiennych stanu i timera, które odczytują lub zapisują wartości, są asynchroniczne. Większość tych operacji zwraca jeden wynik, który pobieramy za pomocą await. Operacje zwracające kolekcję zamiast tego zwracają asynchroniczny iterator, który konsumujesz z .async for Aby poznać wprowadzenie do async/await iteratorów asynchronicznych w Python, zobacz dokumentację Python asyncio.
Poniższa tabela przedstawia operacje zwracające pojedynczy wynik, który można uzyskać za pomocą await:
| Class | Operacje wykorzystujące await |
|---|---|
AsyncValueState |
exists
get, , , updateclear |
AsyncMapState |
exists, getValue, , , updateValue, removeKeycontainsKeyclear |
AsyncListState |
exists, put, , appendValue, appendListclear |
AsyncStatefulProcessorHandle |
registerTimer, deleteTimer |
Poniższa tabela przedstawia operacje zwracające asynchroniczny iterator, który można pobrać za pomocą async for:
| Class | Operacje wykorzystujące async for |
|---|---|
AsyncMapState |
iterator, , keysvalues |
AsyncListState |
get |
AsyncStatefulProcessorHandle |
listTimers |
Przykład: async for
Na przykład, aby odczytać wartości w , AsyncListStateiteruj z takim async for jak w następującym kodzie:
total = 0
async for value in self.items.get():
total += value[0]
Metody tworzące obiekty stanów i usuwające zmienne stanów pozostają synchroniczne: getValueState, getMapState, getListState, oraz deleteIfExists.
Opis każdego typu stanu można znaleźć w artykule Niestandardowe typy stanów.
Optymalizacja z asynchronicznymi wzorcami programowania
Przetwarzanie asynchroniczne jest przydatne, gdy logika czeka na operacje zewnętrzne, takie jak żądania sieciowe. Zamiast czekać na każde żądanie po kolei, używaj asyncio ich jednoczesnego wykonywania i zmniejszania czasu bezczynności.
Przykład: wykonuj żądania współbieżne z asyncio.gather
Poniższy przykład wykorzystuje asyncio.gather jednoczesne uruchamianie wszystkich żądań HTTP na wiersz i czekanie na ich zakończenie, a następnie zapisuje maksymalny wynik w stanie zdarzenia. Zdefiniuj procesor w następującym kodzie:
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()