Richiami di riconoscimento

Un callback di conferma permette al cliente di reagire asincronamente alle conferme e agli errori registrati, senza bloccare il ciclo del produttore. Quando i record diventano durevoli o si guastano, Zerobus Ingest invoca il tuo callback in background, così puoi monitorare i progressi e aggiornare le metriche senza rallentare il produttore, e conoscere i fallimenti non appena accadono.

Questo è diverso da attendere un offset o eseguire lo svuotamento: si tratta di chiamate bloccanti in cui il codice attende in linea che i dati siano resi durevoli. Un callback non è una chiamata bloccante. È un handler che l'SDK invoca per conto tuo quando arrivano le conferme di ricezione.

I callback di conferma sono supportati per flussi SDK JSON e Protocol Buffers (protobuf). I flussi Arrow Flight non supportano i callback; per confermare la durabilità di un flusso Arrow, utilizza wait_for_offset() o flush(). Vedi Usare Arrow Flight con Zerobus Ingest.

I nomi dei metodi e dei tipi qui sotto provengono dall'SDK Python. Altri SDK Zerobus espongono callback di conferma dove supportati, utilizzando costrutti equivalenti in ciascun linguaggio.

Come funzionano i callback

Definisci un callback sottoclassando AckCallback e implementando due metodi:

  • on_ack(offset: int): chiamato quando un invio (un record o un batch) viene confermato con successo come persistente da parte del server. L'offset identifica l'invio confermato.
  • on_error(offset: int, error_message: str): chiamato quando una sottomissione incontra un errore. on_error è facoltativo. Implementalo per gestire o registrare i guasti.

Il callback viene invocato una volta per ogni record o batch inoltrati, quando il relativo offset logico viene confermato oppure quando l'operazione fallisce; rappresenta quindi un indicatore continuo dell'avanzamento dell'ingestione lungo l'intero flusso.

I tuoi metodi di callback girano sui thread in background dell'SDK, quindi invocarli non blocca il tuo produttore. Tienili veloci e non bloccanti. Cosa fare in caso di guasto è responsabilità del tuo cliente: registra, avvisa, riprova o ferma. Alcuni errori sono terminali e, se on_error segnala che il flusso è fallito in modo permanente, devi recuperare su un nuovo stream. Vedi i modelli di recupero e di nuovo tentativo.

Configura una funzione di callback

Alleghi un callback a uno stream passando un'istanza della sottoclasse AckCallback come opzione ack_callback in StreamConfigurationOptions quando crei lo stream. Il callback si applica quindi a ogni record ingerito in quel stream.

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()

Dopo aver registrato il callback, non è necessario attendere in linea. on_ack viene attivato man mano che viene confermata la persistenza di ciascun record e on_error viene attivato se un record non riesce.

Richiami vs. blocchi

Le callback e le chiamate bloccanti risolvono problemi diversi e puoi usarle insieme:

  • Usa una callback di conferma per reagire alle conferme di durabilità e agli errori non appena si verificano, in modo asincrono, senza compromettere un throughput elevato. Ottimo per il monitoraggio dei progressi, le metriche e la registrazione degli errori.
  • Usa wait_for_offset() o flush() quando il codice deve restare bloccato finché un record specifico, o tutti i record in sospeso, non siano stati resi persistenti prima di procedere.