Procesamiento asíncrono con transformWithState (Beta)

Importante

El procesamiento asincrónico para la API basada en filas transformWithState de Python está en Beta. Consulte Azure Databricks versiones de vista previa.

El procesamiento asíncrono está disponible en Databricks Runtime 19 y superiores.

Python transformWithState soporta procesamiento asíncrono basado en asyncio. Al ejecutar operaciones de estado y lógica de usuario simultáneamente a través de claves de agrupación y la comunicación entre procesos en lotes, el procesamiento asíncrono tiene mayor rendimiento que el síncrono, con solo cambios menores en el código. Esta ganancia de rendimiento no requiere ninguna biblioteca asincrónica de terceros. Los usuarios avanzados pueden optimizar aún más sus aplicaciones con patrones de programación asíncronos y bibliotecas habilitadas para asincrónico.

Para usar procesamiento asincrónico, implementa un AsyncStatefulProcessor en lugar del síncrono StatefulProcessor. La AsyncStatefulProcessor API refleja la API síncrona StatefulProcessor , por lo que la mayoría de las aplicaciones requieren solo pequeños cambios para usar la API asíncrona. Véase Implementar un AsyncStatefulProcessorarchivo de .

Para la API síncrona transformWithState y los conceptos centrales, véase Construir una aplicación con estado personalizada con transformWithState.

Note

El procesamiento asíncrono está disponible solo para la API basada en filas de PythontransformWithState. No está soportado ni para transformWithStateInPandas la API de Scala transformWithState . El procesamiento asíncrono no es compatible con la computación sin servidor.

Implementar un AsyncStatefulProcessor

Para convertir un síncrono StatefulProcessor en un AsyncStatefulProcessor, realiza los siguientes cambios:

  • Defina los métodos API (init, close, handleInputRows, handleExpiredTimer, y handleInitialState) con la async def palabra clave.
  • Lee y actualiza los valores de estado y temporizador con await, o ejecutalos usando la biblioteca de asyncio Python. Esto se aplica a operaciones de estado como valueState.get() y a operaciones de temporizador como registerTimer. Crear objetos de estado, como handle.getValueState, permanece síncrono.

Las siguientes consideraciones se aplican al procesamiento asincrónico:

  • Si tu aplicación almacena datos en variables miembros o en sistemas externos, Databricks recomienda reescribir la lógica para que sea seguro para la ejecución concurrente. Dado que handleInputRows y handleExpiredTimer pueden ejecutarse simultáneamente entre claves de agrupación, las ejecuciones entrelazadas no deben corromper los datos compartidos. La mayoría de las solicitudes ya cumplen este requisito.
  • Databricks recomienda no detectar ni suprimir errores de las operaciones de estado. Apache Spark gestiona estos errores por ti. Si una operación de estado falla, Apache Spark falla la tarea y la vuelve a intentar.
    • En un AsyncStatefulProcessor, los errores de operación de estado se gestionan para ti y nunca aparecen en tu código.
    • En un síncrono StatefulProcessor, se generan errores de operación de estado en tu código, pero suprimirlos puede comprometer la corrección de los datos.

Ejemplo: contar filas para cada clave de agrupación

El siguiente ejemplo define un AsyncCountProcessor que cuenta el número de filas de cada clave de agrupación. La value_schema variable define el esquema de el ValueState que almacena el recuento en funcionamiento. En comparación con un síncrono StatefulProcessor, los cambios son la async def palabra clave de cada método y await las operaciones de lectura y actualización de estado. La llamada a getValueState in init permanece sincrónica. Defina el procesador como en el siguiente código:

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

Ejecuta una consulta con un procesador asincrónico

Para ejecutar una consulta con un procesador asincrónico, pasa tu AsyncStatefulProcessor a transformWithState. La consulta utiliza la misma sintaxis que el camino síncrono. Las APIs asíncrona y síncrona comparten el mismo formato de estado, así que puedes cambiar una consulta existente entre una AsyncStatefulProcessor y una síncrona StatefulProcessor mientras reutilizas el mismo punto de control.

Ejemplo: contar eventos en el events conjunto de datos de muestra

El siguiente ejemplo se ejecuta AsyncCountProcessor sobre el events conjunto de datos de muestra. Cada registro tiene un time campo (segundos de época) y un action campo con el valor Open o Close. La consulta agrupa por action y cuenta los eventos para cada tipo de acción. Para más conjuntos de datos de muestra, véase Conjuntos de datos de muestra.

La input_schema variable define el esquema de los registros fuente, y la output_schema variable define el esquema de las filas que emite el procesador. Para leer el conjunto de datos de muestra como un flujo, define ambos esquemas y luego inicia la consulta como en el siguiente código:

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

Una vez completada la consulta, consulta el recuento de acción para cada tipo de acción tal como en el siguiente código:

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

Estado asíncrono y operaciones del temporizador

En un AsyncStatefulProcessor, las operaciones de variables de estado y temporizador que leen o escriben valores son asíncronas. La mayoría de estas operaciones devolven un único resultado que recuperas con await. Las operaciones que devuelven una colección en su lugar devuelven un iterador asíncrono que consumes con async for. Para una introducción a async/await los iteradores asíncronos en Python, consulta la documentación de asíncro en Python.

La siguiente tabla enumera las operaciones que devuelven un único resultado que puedes recuperar con await:

Class Operaciones que utilizan await
AsyncValueState exists, get, , update, clear
AsyncMapState exists, getValue, containsKey, updateValue, removeKey, clear
AsyncListState exists, put, appendValue, appendList, clear
AsyncStatefulProcessorHandle registerTimer, deleteTimer

La siguiente tabla lista las operaciones que devuelven un iterador asincrónico que puedes recuperar con async for:

Class Operaciones que utilizan async for
AsyncMapState iterator, keys, values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

Ejemplo: async for

Por ejemplo, para leer los valores en un AsyncListState, itera con async for como en el siguiente código:

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

Los métodos que crean objetos de estado y eliminan variables de estado permanecen síncronos: getValueState, getMapState, getListState, y deleteIfExists.

Para una descripción de cada tipo de estado, véase Tipos de estado personalizados.

Optimizar con patrones de programación asincrónicos

El procesamiento asíncrono es útil cuando tu lógica espera operaciones externas, como solicitudes de red. En lugar de esperar cada solicitud en secuencia, úsalo asyncio para ejecutarlas simultáneamente y reducir el tiempo de inactividad.

Ejemplo: ejecutar solicitudes concurrentes con asyncio.gather

El siguiente ejemplo se utiliza asyncio.gather para disparar todas las solicitudes HTTP por fila simultáneamente y esperar a que se completen, para luego almacenar la puntuación máxima en el estado. Defina el procesador como en el siguiente código:

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