Callbacks de reconhecimento
Um callback de confirmação permite que seu cliente reaja a confirmações de registro e erros de forma assíncrona, sem bloquear seu loop de produtor. À medida que os registros se tornam duráveis ou falham, o Zerobus Ingest invoca seu callback em segundo plano, para que você possa acompanhar o progresso e atualizar as métricas sem reduzir a velocidade do seu produtor, e saber sobre as falhas assim que elas ocorrerem.
Isso difere de aguardar por um offset ou flushing: essas são chamadas de bloqueio onde seu código aguarda pela durabilidade em linha. Um callback não é uma chamada de bloqueio. É um manipulador que o SDK invoca para você quando as confirmações chegam.
Callbacks de confirmação são suportados para transmissões de SDK JSON e Protocol Buffers (protobuf). As transmissões do Arrow Flight não oferecem suporte a callbacks; para confirmar a durabilidade em uma transmissão Arrow, use wait_for_offset() ou flush(). Veja Usar o Arrow Flight com o Zerobus Ingest.
Os nomes de métodos e tipos abaixo são do SDK para Python. Outros SDKs do Zerobus expõem callbacks de reconhecimento onde houver suporte, usando construções equivalentes em cada linguagem.
Como funcionam os callbacks
Você define um callback criando uma subclasse de AckCallback e implementando dois métodos:
on_ack(offset: int): chamado quando uma submissão (um registro ou um lote) é reconhecida com sucesso como durável pelo servidor. Ooffsetidentifica a submissão reconhecida.on_error(offset: int, error_message: str): chamado quando um envio encontra um erro.on_erroré opcional. Implemente-o para tratar ou registrar falhas.
O callback é invocado uma vez para cada registro ou lote enviado quando seu deslocamento lógico é reconhecido ou falha, portanto, é um sinal contínuo do progresso da ingestão em toda a transmissão.
Seus métodos de callback são executados nos threads em segundo plano do SDK, portanto, invocá-los não bloqueia seu produtor. Mantenha-os rápidos e sem bloqueios. O que fazer em caso de falha é responsabilidade do seu cliente: registrar em logs, emitir alerta, repetir ou parar. Alguns erros são terminais e, se on_error informar que a transmissão falhou permanentemente, você deve recuperar em uma nova transmissão. Consulte Padrões de recuperação e repetição.
Configurar um callback
Você anexa um callback a uma transmissão passando uma instância da sua subclasse AckCallback como a opção ack_callback em StreamConfigurationOptions ao criar a transmissão. O callback então se aplica a cada registro ingerido nessa transmissão.
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import AckCallback, StreamConfigurationOptions, TableProperties
class MyAckCallback(AckCallback):
def on_ack(self, offset: int) -> None:
print(f"Record acknowledged at offset: {offset}")
def on_error(self, offset: int, error_message: str) -> None:
print(f"Error at offset {offset}: {error_message}")
sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)
table_properties = TableProperties("main.default.air_quality")
options = StreamConfigurationOptions(
ack_callback=MyAckCallback(),
)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties, options)
try:
for row in records:
stream.ingest_record_offset(row)
finally:
stream.close()
Com o callback registrado, você não precisa aguardar inline. on_ack é disparado à medida que cada registro é confirmado como durável, e on_error é disparado se um registro falhar.
Callbacks vs. bloqueio
Callbacks e chamadas de bloqueio resolvem problemas diferentes, e você pode usá-los juntos:
- Use um callback de confirmação para reagir a confirmações de durabilidade e erros conforme eles ocorrem, de forma assíncrona, enquanto mantém um throughput alto. Adequado para acompanhamento de progresso, métricas e registro de erros.
- Use
wait_for_offset()ouflush()quando seu código precisar bloquear até que um registro específico, ou todos os registros pendentes, estejam duráveis antes de prosseguir.
Relacionado
- Bloqueio e confirmação de mensagens: bloqueio na durabilidade com
wait_for_offseteflush. - Padrões de recuperação e nova tentativa: Tratamento de erros e recuperação de registros não reconhecidos.
- Tratamento de erros do Zerobus Ingest: referência de código de erro.