Bevestigingsterugroepen

Een acquiledgment callback stelt je klant in staat om asynchroon te reageren op recordbevestigingen en fouten, zonder je producer-lus te blokkeren. Naarmate records duurzaam worden of falen, activeert Zerobus Ingest je callback op de achtergrond, zodat je de voortgang kunt volgen en statistieken kunt bijwerken zonder je producer te vertragen, en direct over fouten kunt leren zodra ze optreden.

Dit verschilt van wachten op een offset of flushen: dat zijn blokkerende aanroepen waarbij je code direct wacht op duurzame opslag. Een callback is geen blokkeringsoproep. Het is een handler die de SDK voor je aanroept wanneer bevestigingen binnenkomen.

Bevestigingscallbacks worden ondersteund voor JSON- en Protocol Buffers (protobuf) SDK-stromen. Arrow Flight-streams ondersteunen geen callbacks; om de duurzaamheid van een Arrow-stream te bevestigen, gebruik wait_for_offset() of flush(). Zie Arrow Flight gebruiken met Zerobus Ingest.

Methode- en typenamen hieronder komen uit de Python SDK. Andere Zerobus SDK's bieden bevestigingscallbacks waar ondersteund, met equivalente constructies in elke taal.

Hoe aanroepfuncties werken

Je definieert een callback door AckCallback te subklasseren en twee methoden te implementeren:

  • on_ack(offset: int): wordt aangeroepen wanneer een inzending (een record of batch) succesvol als duurzaam wordt bevestigd door de server. De offset identificeert de erkende inzending.
  • on_error(offset: int, error_message: str): wordt opgeroepen wanneer een inzending een fout ondervindt. on_error is een optie. Implementeer het om storingen te verwerken of te loggen.

De callback wordt één keer aangeroepen voor elk verzonden record of elke batch zodra de logische offset wordt bevestigd of de verwerking mislukt, en vormt daarmee een doorlopend signaal van de voortgang van de inname binnen de stream.

Je callback-methoden draaien op de achtergrondthreads van de SDK, dus het aanroepen ervan blokkeert je producent niet. Zorg dat ze snel zijn en niet blokkeren. Wat je doet bij een mislukking is de verantwoordelijkheid van je klant: loggen, waarschuwen, opnieuw proberen of stoppen. Sommige fouten zijn fataal, en als on_error meldt dat de stream permanent heeft gefaald, moet je herstellen met een nieuwe stream. Zie Herstel- en herpogingspatronen.

Configureer een callback

Je koppelt een callback aan een stream door bij het aanmaken van de stream een instantie van je AckCallback-subklasse mee te geven als de ack_callback-optie in StreamConfigurationOptions. De callback geldt dan voor elk record dat in die stream wordt opgenomen.

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

Zodra het terugbelverzoek is geregistreerd, hoef je niet meer aan de lijn te wachten. on_ack brandt wanneer elk record als duurzaam wordt bevestigd, en on_error brandt als een record faalt.

Callbacks versus blokkeren

Callbacks en blokkeringsaanroepen lossen verschillende problemen op, en je kunt ze samen gebruiken:

  • Gebruik een bevestigingscallback om te reageren op duurzaamheidsbevestigingen en fouten zodra ze optreden, asynchroon, terwijl de doorvoer hoog blijft. Goed voor voortgangstracking, statistieken en foutregistratie.
  • Gebruik wait_for_offset() of flush() wanneer je code moet blokkeren totdat een specifiek record, of alle lopende records, duurzaam zijn voordat je verder gaat.