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​