Asynchrone Verarbeitung mit transformWithState (Beta)

Important

Die asynchrone Verarbeitung für die zeilenbasierte transformWithState API von Python befindet sich in der Beta-Phase. Siehe Azure Databricks Preview-Versionen.

Asynchrone Verarbeitung ist ab Databricks Runtime 19 verfügbar.

Python transformWithState unterstützt asynchrone Verarbeitung, die auf asynciobasiert. Durch das gleichzeitige Ausführen von Zustandsoperationen und Benutzerlogik über Schlüssel hinweg und das Batching der Kommunikation zwischen Prozessen hat asynchrone Verarbeitung einen höheren Durchsatz als synchrone Verarbeitung mit nur geringfügigen Codeänderungen. Dieser Durchsatzgewinn erfordert keine asynchronen Drittanbieter-Bibliotheken. Fortgeschrittene Nutzer können ihre Anwendungen weiter mit asynchronen Programmiermustern und asynk-fähigen Bibliotheken optimieren.

Um asynchrone Verarbeitung zu verwenden, implementieren Sie anstelle AsyncStatefulProcessor des synchronen StatefulProcessor. Die AsyncStatefulProcessor API spiegelt die synchrone StatefulProcessor API wider, sodass die meisten Anwendungen nur geringe Änderungen benötigen, um die asynchrone API zu nutzen. Siehe Implementierung eines AsyncStatefulProcessor.

Für die synchrone transformWithState API und Kernkonzepte siehe Build a custom stateful application with transformWithState.

Note

Asynchrone Verarbeitung ist nur für die zeilenbasierte transformWithState API von Python verfügbar. Sie wird weder für transformWithStateInPandas noch für die Scala-API transformWithState unterstützt. Asynchrone Verarbeitung wird in serverloser Berechnung nicht unterstützt.

Implementiere ein AsyncStatefulProcessor

Um eine synchrone StatefulProcessor in eine umzuwandeln AsyncStatefulProcessor, nehmen Sie folgende Änderungen vor:

  • Definiere die API-Methoden (init, close, handleInputRows, handleExpiredTimer, und handleInitialState) mit dem async def Schlüsselwort.
  • Lesen und aktualisieren Sie Zustands- und Timerwerte mit await, oder führen Sie sie mit der Python-Bibliothek asyncio aus. Dies gilt für Zustandsoperationen wie valueState.get() und für Timer-Operationen registerTimerwie . Das Erstellen von Zustandsobjekten, wie handle.getValueState, bleibt synchron.

Die folgenden Überlegungen gelten für asynchrone Verarbeitung:

  • Wenn Ihre Anwendung Daten in Mitgliedsvariablen oder in externen Systemen speichert, empfiehlt Databricks, die Logik so umzuschreiben, dass sie für gleichzeitige Ausführung sicher ist. Da handleInputRows und handleExpiredTimer gleichzeitig über Gruppierungsschlüssel hinweg ausgeführt werden können, dürfen interleaved-Läufe die gemeinsamen Daten nicht beschädigen. Die meisten Bewerbungen erfüllen diese Anforderung bereits.
  • Databricks empfiehlt, Fehler bei Zustandsoperationen nicht zu erkennen oder zu unterdrücken. Apache Spark übernimmt diese Fehler für Sie. Wenn eine Zustandsoperation fehlschlägt, scheitert Apache Spark an der Aufgabe und versucht es erneut.
    • In einem AsyncStatefulProcessorZustand werden Betriebsfehler für Sie verwaltet und tauchen nie in Ihrem Code auf.
    • In einem synchronen StatefulProcessorZustand werden Betriebsfehler in Ihrem Code angezeigt, aber deren Unterdrückung kann die Datenkorrektheit beeinträchtigen.

Beispiel: Zählzeilen für jede Gruppierungstaste

Das folgende Beispiel definiert ein, AsyncCountProcessor der die Anzahl der Zeilen für jeden Gruppierungsschlüssel zählt. Die Variable value_schema definiert das Schema von , ValueState das die laufende Anzahl speichert. Im Vergleich zu einer synchronen StatefulProcessorMethode sind die Änderungen das async def Schlüsselwort bei jeder Methode sowie await bei den Zustandslese- und Aktualisierungsoperationen. Der Anruf nach getValueState innen init bleibt synchron. Definieren Sie den Prozessor wie im folgenden 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

Führe eine Abfrage mit einem asynchronen Prozessor aus

Um eine Abfrage mit einem asynchronen Prozessor auszuführen, übergebe dein AsyncStatefulProcessor an transformWithState. Die Abfrage verwendet dieselbe Syntax wie der synchrone Pfad. Die asynchronen und synchronen APIs teilen dasselbe Zustandsformat, sodass man eine bestehende Abfrage zwischen einer AsyncStatefulProcessor und einer synchronen StatefulProcessor Stelle wechseln kann, während derselbe Checkpoint wiederverwendet wird.

Beispiel: Zähle Ereignisse im events Beispieldatensatz

Das folgende Beispiel läuft AsyncCountProcessor mit dem events Beispieldatensatz. Jeder Datensatz hat ein Feld time (Epochensekunden) und ein action Feld mit dem Wert Open oder Close. Die Abfrage gruppiert nach action und zählt die Ereignisse für jeden Aktionstyp. Weitere Beispieldatensätze finden Sie unter Beispieldatensätze.

Die Variable input_schema definiert das Schema der Quelldatensätze, und die Variable output_schema definiert das Schema der Zeilen, die der Prozessor ausgibt. Um den Beispieldatensatz als Strom auszulesen, definieren Sie beide Schemata und starten dann die Abfrage wie im folgenden 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()

Nachdem die Abfrage abgeschlossen ist, sehen Sie sich die laufende Anzahl für jeden Aktionstyp wie im folgenden Code an:

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

Asynchrone Zustands- und Timeroperationen

In einem AsyncStatefulProcessorZustand sind Variablen- und Timeroperationen, die Werte lesen oder schreiben, asynchron. Die meisten dieser Operationen liefern ein einziges Ergebnis, das Sie mit awaitabrufen. Operationen, die eine Sammlung zurückgeben, geben stattdessen einen asynchronen Iterator zurück, den Sie mit async forkonsumieren. Eine Einführung in async/await und asynchrone Iteratoren in Python finden Sie in der Python asyncio-Dokumentation.

Die folgende Tabelle listet Operationen auf, die ein einzelnes Ergebnis zurückgeben, das Sie mit awaitabrufen können:

Class Operationen, die verwenden await
AsyncValueState exists, get, update, clear
AsyncMapState exists, getValue, containsKey, updateValue, , removeKey, clear
AsyncListState exists, put, appendValue, appendList, , clear
AsyncStatefulProcessorHandle registerTimer, deleteTimer

Die folgende Tabelle listet Operationen auf, die einen asynchronen Iterator zurückgeben, mit dem Sie abrufen async forkönnen:

Class Operationen, die verwenden async for
AsyncMapState iterator, keys, values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

Beispiel: async for

Zum Beispiel, um die Werte in einem zu lesen AsyncListState, iterieren Sie mit async for wie im folgenden Code:

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

Die Methoden, die Zustandsobjekte erzeugen und Zustandsvariablen löschen, bleiben synchron: getValueState, getMapState, , getListStateund deleteIfExists.

Für eine Beschreibung jedes Zustandstyps siehe Benutzerdefinierte Zustandstypen.

Optimieren Sie mit asynchronen Programmierungsmustern

Asynchrone Verarbeitung ist nützlich, wenn Ihre Logik auf externe Operationen wie Netzwerkanfragen wartet. Anstatt auf jede Anfrage nacheinander zu warten, sollten asyncio Sie die Anfragen gleichzeitig ausführen und die Leerlaufzeit reduzieren.

Beispiel: Gleiche Anfragen ausführen mit asyncio.gather

Das folgende Beispiel verwendet asyncio.gather , um alle HTTP-Anfragen pro Zeile gleichzeitig auszulösen und auf deren Abschluss zu warten, dann speichert es den maximalen Score im Zustand. Definieren Sie den Prozessor wie im folgenden 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()