Processamento assíncrono com transformWithState (Beta)

Important

O processamento assíncrono para a API baseada em linhas transformWithState em Python está em Beta. Consulte as versões de pré-visualização do Azure Databricks.

O processamento assíncrono está disponível no Databricks Runtime 19 e superiores.

Python transformWithState suporta processamento assíncrono baseado em asyncio. Ao executar operações de estado e lógica de utilizador em simultâneo através de chaves de agrupamento e comunicação em lote entre processos, o processamento assíncrono tem maior rendimento do que o processamento síncrono, com apenas pequenas alterações de código. Este ganho de rendimento não requer quaisquer bibliotecas assíncronas de terceiros. Utilizadores avançados podem otimizar ainda mais as suas aplicações com padrões de programação assíncronos e bibliotecas habilitadas para assíncrono.

Para usar processamento assíncrono, implemente um AsyncStatefulProcessor em vez do síncrono StatefulProcessor. A AsyncStatefulProcessor API espelha a API síncrona StatefulProcessor , pelo que a maioria das aplicações requer apenas pequenas alterações para usar a API assíncrona. Ver implementar um AsyncStatefulProcessor.

Para a API síncrona transformWithState e os conceitos centrais, veja Construir uma aplicação com estado personalizada com transformWithState.

Note

O processamento assíncrono está disponível apenas para a API baseada em linhas transformWithState em Python. Não é suportado nem para transformWithStateInPandas a API Scala transformWithState . O processamento assíncrono não é suportado em computação serverless.

Implementar um AsyncStatefulProcessor

Para converter um síncrono StatefulProcessor para um AsyncStatefulProcessor, faça as seguintes alterações:

  • Defina os métodos API (init, close, handleInputRows, handleExpiredTimer, e handleInitialState) com a async def palavra-chave.
  • Leia e atualize os valores de estado e temporizador com await, ou execute-os usando a biblioteca do asyncio Python. Isto aplica-se a operações de estado como valueState.get() e a operações de temporizador como registerTimer. Criar objetos de estado, como handle.getValueState, mantém-se síncrono.

As seguintes considerações aplicam-se ao processamento assíncrono:

  • Se a sua aplicação armazenar dados em variáveis membros ou em sistemas externos, o Databricks recomenda que reescreva a lógica para ser seguro para execução simultânea. Como handleInputRows e handleExpiredTimer podem correr simultaneamente entre chaves de agrupamento, as execuções entrelaçadas não devem corromper dados partilhados. A maioria das candidaturas já cumpre este requisito.
  • O Databricks recomenda que não apanhe nem suprima erros das operações de estado. O Apache Spark trata destes erros por si. Se uma operação de estado falhar, o Apache Spark falha a tarefa e tenta novamente.
    • Num AsyncStatefulProcessor, erros de operação de estado são geridos para si e nunca surgem no seu código.
    • Num síncrono StatefulProcessor, erros de operação de estado surgem no seu código, mas suprimi-los pode comprometer a correção dos dados.

Exemplo: contar as linhas para cada chave de agrupamento

O exemplo seguinte define um AsyncCountProcessor que conta o número de linhas para cada chave de agrupamento. A value_schema variável define o esquema do ValueState que armazena a contagem em curso. Comparado com um síncrono StatefulProcessor, as alterações são a async def palavra-chave em cada método e await nas operações de leitura e atualização de estado. A chamada para getValueState in init mantém-se síncrona. Defina o processador como no seguinte 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

Execute uma consulta com um processador assíncrono

Para executar uma consulta com um processador assíncrono, passe your AsyncStatefulProcessor para transformWithState. A consulta utiliza a mesma sintaxe do caminho síncrono. As APIs assíncrona e síncrona partilham o mesmo formato de estado, por isso pode alternar uma consulta existente entre uma AsyncStatefulProcessor e uma síncrona StatefulProcessor enquanto reutiliza o mesmo checkpoint.

Exemplo: contar eventos no events conjunto de dados de amostra

O exemplo seguinte corre AsyncCountProcessor contra o events conjunto de dados de amostra. Cada registo tem um time campo (segundos de época) e um action campo com o valor Open ou Close. A consulta agrupa por action e conta os eventos para cada tipo de ação. Para mais conjuntos de dados de exemplo, consulte Conjuntos de dados de exemplo.

A input_schema variável define o esquema dos registos fonte, e a output_schema variável define o esquema das linhas que o processador emite. Para ler o conjunto de dados de exemplo como um fluxo, defina ambos os esquemas e depois inicie a consulta conforme no seguinte 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()

Após a conclusão da consulta, veja a contagem contínua para cada tipo de ação conforme no seguinte código:

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

Estado assíncrono e operações do temporizador

Num AsyncStatefulProcessor, as operações de variável de estado e temporizador que leem ou escrevem valores são assíncronas. A maioria destas operações retorna um único resultado que recupera com await. Operações que retornam uma coleção retornam em vez disso um iterador assíncrono que consome com async for. Para uma introdução a async/await iteradores assíncronos em Python, consulte a documentação sobre assíncronos em Python.

A tabela seguinte lista operações que retornam um único resultado que pode obter com await:

Class Operações que utilizam await
AsyncValueState exists, get, update, clear
AsyncMapState exists, getValue, containsKey, updateValue, removeKey, clear
AsyncListState exists, put, appendValue, appendList, clear
AsyncStatefulProcessorHandle registerTimer, deleteTimer

A tabela seguinte lista as operações que retornam um iterador assíncrono que pode obter com async for:

Class Operações que utilizam async for
AsyncMapState iterator, keys, values
AsyncListState get
AsyncStatefulProcessorHandle listTimers

Exemplo: async for

Por exemplo, para ler os valores num AsyncListState, itere com async for como no seguinte código:

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

Os métodos que criam objetos de estado e eliminam variáveis de estado mantêm-se síncronos: getValueState, getMapState, getListState, e deleteIfExists.

Para uma descrição de cada tipo de estado, veja Tipos de estado personalizados.

Otimizar com padrões de programação assíncronos

O processamento assíncrono é útil quando a sua lógica espera em operações externas, como pedidos de rede. Em vez de esperar por cada pedido em sequência, use-se asyncio para executar os pedidos em simultâneo e reduzir o tempo de inatividade.

Exemplo: executar pedidos concorrentes com asyncio.gather

O exemplo seguinte serve asyncio.gather para disparar todos os pedidos HTTP por linha em simultâneo e esperar que sejam concluídos, armazenando depois a pontuação máxima no estado. Defina o processador como no seguinte 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()