Pular para o conteúdo principal

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. O offset identifica 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.

Python
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() ou flush() quando seu código precisar bloquear até que um registro específico, ou todos os registros pendentes, estejam duráveis antes de prosseguir.

Relacionado