Processamento assíncrono com transformWithState (Beta)
Beta
O processamento assíncrono para a API transformWithState baseada em linhas do Python está em versão Beta. Consulte Lançamentos de pré-visualização do Databricks.
O processamento assíncrono está disponível no Databricks Runtime 19 e acima.
O Python transformWithState oferece suporte ao processamento assíncrono criado com base em asyncio. Ao executar operações de estado e lógica de usuário simultaneamente entre chaves de agrupamento e processar em lotes a comunicação entre processos, o processamento assíncrono tem um throughput maior do que o processamento síncrono com apenas pequenas alterações no código. Esse ganho de throughput não requer nenhuma biblioteca assíncrona de terceiros. Usuários avançados podem otimizar ainda mais suas aplicações com padrões de programação assíncrona e bibliotecas habilitadas para assincronia.
Para usar o processamento assíncrono, implemente um AsyncStatefulProcessor em vez do StatefulProcessor síncrono. A API AsyncStatefulProcessor espelha a API StatefulProcessor síncrona, portanto, a maioria dos aplicativos requer apenas pequenas alterações para usar a API assíncrona. Consulte Implementar um AsyncStatefulProcessor.
Para a API transformWithState síncrona e conceitos fundamentais, consulte Criar um aplicativo com estado personalizado com transformWithState.
O processamento assíncrono está disponível apenas para a API transformWithState baseada em linhas do Python. Não é compatível com transformWithStateInPandas ou com a API transformWithState do Scala. O processamento assíncrono não é compatível com compute serverless.
Implementar um AsyncStatefulProcessor
Para converter um StatefulProcessor síncrono em um AsyncStatefulProcessor, faça as seguintes alterações:
- Defina os métodos da API (
init,close,handleInputRows,handleExpiredTimerehandleInitialState) com a palavra-chaveasync def. - Leia e atualize valores de estado e de temporizador com
await, ou execute-os usando a bibliotecaasynciodo Python. Isso se aplica a operações de estado, comovalueState.get(), e a operações de temporizador, comoregisterTimer. A criação de objetos de estado, comohandle.getValueState, permanece síncrona.
As seguintes considerações se aplicam ao processamento assíncrono:
- Se sua aplicação armazena dados em variáveis de membro ou em sistemas externos, o Databricks recomenda que você reescreva a lógica para que ela seja segura para execução concorrente. Como
handleInputRowsehandleExpiredTimerpodem ser executados de forma concorrente entre chaves de agrupamento, execuções intercaladas não devem corromper dados compartilhados. A maioria das aplicações já atende a esse requisito. - A Databricks recomenda que você não capture ou suprima erros de operações de estado. O Apache Spark trata esses erros para você. Se uma operação de estado falhar, o Apache Spark falhará a tarefa e a reexecutará.
- Em um
AsyncStatefulProcessor, os erros de operação de estado são gerenciados para você e nunca são exibidos para o seu código. - Em um
StatefulProcessorsíncrono, erros de operações de estado são gerados em seu código, mas suprimi-los pode comprometer a correção dos dados comprometidos.
- Em um
Exemplo: contar linhas para cada key de agrupamento
O exemplo a seguir define um AsyncCountProcessor que conta o número de linhas para cada key de agrupamento. A variável value_schema define o esquema do ValueState que armazena a contagem em execução. Comparado a um StatefulProcessor síncrono, as alterações são a palavra-chave async def em cada método e await nas operações de leitura e atualização de estado. A chamada para getValueState em init permanece síncrona. Defina o processador como no código a seguir:
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
Executar uma query com um processador assíncrono
Para executar uma query com um processador assíncrono, passe seu AsyncStatefulProcessor para transformWithState. A query usa a mesma sintaxe que o caminho síncrono. As APIs assíncronas e síncronas compartilham o mesmo formato de estado, portanto, você pode alternar uma query existente entre uma AsyncStatefulProcessor e uma StatefulProcessor síncrona enquanto reutiliza o mesmo ponto de verificação.
Exemplo: contar eventos no dataset de amostra events
O exemplo a seguir executa AsyncCountProcessor no dataset de amostra events. Cada registro tem um campo time (segundos da época) e um campo action com o valor Open ou Close. A query agrupa por action e conta os eventos para cada tipo de ação. Para mais datasets de amostra, consulte Datasets de exemplo.
A variável input_schema define o esquema dos registros de origem, e a variável output_schema define o esquema das linhas que o processador emite. Para ler o dataset de exemplo como uma transmissão, defina ambos os esquemas e, em seguida, inicie a query como no código a seguir:
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 query, view a contagem em execução para cada tipo de ação conforme o código a seguir:
display(spark.sql("SELECT action, MAX(count) AS count FROM async_counts GROUP BY action ORDER BY action"))
Operações de estado assíncrono e timer
Em um AsyncStatefulProcessor, as operações de variável de estado e de temporizador que leem ou gravam valores são assíncronas. A maioria dessas operações retorna um único resultado que você recupera com await. As operações que retornam uma coleção retornam, em vez disso, um iterador assíncrono que você consome com async for. Para uma introdução ao async/await e iteradores assíncronos em Python, consulte a documentação do asyncio do Python.
A tabela a seguir lista as operações que retornam um único resultado que você pode recuperar com await:
Aula | Operações que usam |
|---|---|
|
|
|
|
|
|
|
|
A tabela a seguir lista as operações que retornam um iterador assíncrono que você pode recuperar com async for:
Aula | Operações que usam |
|---|---|
|
|
|
|
|
|
Exemplo: async for
Por exemplo, para ler os valores em um AsyncListState, itere com async for como no código a seguir:
total = 0
async for value in self.items.get():
total += value[0]
Os métodos que criam objetos de estado e excluem variáveis de estado permanecem síncronos: getValueState, getMapState, getListState e deleteIfExists.
Para uma descrição de cada tipo de estado, consulte Tipos de estado personalizados.
Otimize com padrões de programação assíncrona
O processamento assíncrono é útil quando sua lógica aguarda operações externas, como solicitações de rede. Em vez de aguardar cada solicitação em sequência, use asyncio para executar as solicitações simultaneamente e reduzir o tempo parado.
Exemplo: executar solicitações concorrentes com asyncio.gather
O exemplo a seguir usa asyncio.gather para disparar todas as solicitações HTTP por linha simultaneamente e aguardar a conclusão delas, e então armazena a pontuação máxima no estado. Defina o processador conforme o código a seguir:
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()