Retornos de chamada de reconhecimento

Um retorno de chamada de reconhecimento permite que seu cliente reaja aos reconhecimentos e erros de gravação de forma assíncrona, sem bloquear o loop do produtor. À medida que os registros se tornam duráveis ou falham, o Zerobus Ingest invoca seu retorno de chamada em segundo plano para que você possa acompanhar o progresso e atualizar métricas sem atrasar o produtor, além de aprender sobre falhas assim que acontecem.

Isso difere de esperar por um deslocamento ou liberação: essas são chamadas de bloqueio onde seu código espera pela durabilidade em linha. Um retorno de chamada não é uma chamada de bloqueio. É um manipulador que o SDK chama para você quando as confirmações chegam.

Retornos de chamada de reconhecimento têm suporte para fluxos de SDK JSON e Buffers de Protocolo (protobuf). Os fluxos do Arrow Flight não dão suporte a retornos de chamada; para confirmar a durabilidade em um fluxo do Arrow, use wait_for_offset() ou flush(). Consulte Usar Arrow Flight com Zerobus Ingest.

Os nomes dos métodos e tipos abaixo são do SDK Python. Outros SDKs do Zerobus expõem retornos de chamada de reconhecimento, quando compatíveis, usando mecanismos equivalentes em cada linguagem.

Como funcionam os callbacks

Você define um callback subclassificando AckCallback e implementando dois métodos:

  • on_ack(offset: int): chamada quando uma submissão (um registro ou lote) é reconhecida com sucesso como durável pelo servidor. O offset identifica o envio confirmado.
  • on_error(offset: int, error_message: str): chamada quando uma submissão encontra um erro. on_error é opcional. Implemente-o para lidar com falhas ou registrá-las.

O retorno de chamada é invocado uma vez para cada registro ou lote enviado quando o deslocamento lógico é confirmado ou falha; portanto, ele é um indicador contínuo do progresso da ingestão ao longo do fluxo.

Seus métodos de retorno de chamada são executados nas threads de 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, alertar, tentar novamente ou parar. Alguns erros são terminais; se on_error relatar que a transmissão falhou permanentemente, você deverá recuperar em uma nova transmissão. Veja Padrões de recuperação e retentativas.

Configurar um retorno de chamada

Você anexa um retorno de chamada a um fluxo passando uma instância da sua subclasse AckCallback como a opção ack_callback em StreamConfigurationOptions quando cria o fluxo. O retorno de chamada se aplica a cada registro ingerido nesse fluxo.

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 retorno de chamada registrado, você não aguarda na fila. on_ack é disparado à medida que a durabilidade de cada registro é confirmada, e on_error é disparado se um registro falha.

Retornos de chamada x bloqueio

Retornos de chamada e chamadas bloqueadas resolvem diferentes problemas. Você pode usá-los juntos:

  • Use um retorno de chamada de reconhecimento para reagir a confirmações de durabilidade e erros à medida que ocorrem, de forma assíncrona, mantendo a alta taxa de transferência. Bom para acompanhamento de progresso, métricas e registro de erros.
  • Use wait_for_offset() ou flush() quando seu código precisar ser bloqueado até que um registro específico, ou todos os registros pendentes, sejam duráveis antes de prosseguir.