Retornos de reconhecimento

Um callback de confirmação permite que o seu cliente reaja de forma assíncrona a confirmações e erros de registo, sem bloquear o seu loop de produtor. À medida que os registos se tornam duráveis ou falham, o Zerobus Ingest invoca o seu callback em segundo plano, para que possa acompanhar o progresso e atualizar métricas sem atrasar o seu produtor, e aprender sobre falhas assim que acontecem.

Isto difere de esperar por um offset ou flushing: são chamadas de bloqueio onde o teu código espera pela durabilidade em linha. Uma chamada de retorno não é uma chamada de bloqueio. É um processador que o SDK invoca por si quando as confirmações chegam.

Os callbacks de confirmação são suportados para fluxos de SDK JSON e Protocol Buffers (protobuf). Os fluxos do Arrow Flight não suportam callbacks; para confirmar a durabilidade num fluxo Arrow, use wait_for_offset() ou flush(). Veja 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 funções de retorno de confirmação, quando tal é suportado, usando mecanismos equivalentes em cada linguagem.

Como funcionam os callbacks

Defines um callback subclassificando AckCallback e implementando dois métodos:

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

O callback é invocado uma vez por cada registo ou lote submetido quando o respetivo offset lógico é confirmado ou ocorre uma falha, pelo que constitui um sinal contínuo do progresso da ingestão em todo o fluxo.

Os teus métodos de callback correm nos threads de fundo do SDK, por isso invocá-los não bloqueia o teu produtor. Mantenha-os rápidos e não bloqueantes. O que fazer em relação a uma falha é responsabilidade do seu cliente: registar, alertar, tentar novamente ou parar. Alguns erros são terminais e, se on_error indicar que o fluxo falhou permanentemente, tem de recuperar num novo fluxo. Veja Padrões de recuperação e repetição.

Configurar um retorno de chamada

Anexa um callback a um stream, passando uma instância da tua subclasse AckCallback como a opção ack_callback em StreamConfigurationOptions ao criares o stream. O callback aplica-se então a todos os registos ingeridos 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 callback registado, não precisa de esperar em linha. on_ack é acionado à medida que cada registo é confirmado como persistente, e on_error é acionado se um registo falhar.

Chamadas de retorno vs. bloqueio

As funções de retorno e as chamadas bloqueantes resolvem problemas diferentes, e podem ser utilizadas em conjunto:

  • Use um callback de confirmação para reagir às confirmações e erros de durabilidade à medida que acontecem, de forma assíncrona, mantendo um alto rendimento. Bom para acompanhamento de progresso, métricas e registo de erros.
  • Use wait_for_offset() ou flush() quando o seu código tiver de ser bloqueado até que um registo específico, ou todos os registos pendentes, sejam duráveis antes de prosseguir.