Rückrufe zur Bestätigung

Mit einem Anerkennungsrückruf kann Ihr Client, asynchron auf Anerkennungen und Fehler reagieren, ohne Ihre Producerschleife zu blockieren. Wenn Datensätze dauerhaft werden oder fehlschlagen, ruft Zerobus Ingest Ihren Rückruf im Hintergrund auf, damit Sie den Fortschritt verfolgen und Metriken aktualisieren können, ohne Ihren Producer zu verlangsamen, und über Fehler informiert werden, sobald sie auftreten.

Dies unterscheidet sich vom Warten auf einen Offset oder Flushing: Diese blockieren Aufrufe, bei denen Ihr Code inline auf Dauerhaftigkeit wartet. Ein Rückruf ist kein Blockanruf. Es ist ein Handler, den das SDK für dich aufruft, wenn Bestätigungen eintreffen.

Bestätigungsrückrufe werden für JSON- und Protocol Buffers (protobuf) SDK-Streams unterstützt. Arrow Flight-Streams unterstützen keine Callbacks; um die Dauerhaftigkeit bei einem Arrow-Stream zu bestätigen, verwenden Sie wait_for_offset() oder flush(). Siehe Verwendung von Arrow Flight mit Zerobus Ingest.

Die unten aufgeführten Methoden- und Typnamen stammen aus dem Python SDK. Andere Zerobus-SDKs stellen Bestätigungsrückrufe bereit, wo sie unterstützt werden, und verwenden entsprechende Konstrukte in jeder Sprache.

Wie Rückrufe funktionieren

Sie definieren einen Rückruf, indem Sie zwei Methoden in Unterklassen unterteilen AckCallback und implementieren:

  • on_ack(offset: int): aufgerufen, wenn eine Einreichung (ein Datensatz oder ein Batch) vom Server erfolgreich als dauerhaft bestätigt wird. Das offset kennzeichnet die bestätigte Übermittlung.
  • on_error(offset: int, error_message: str): aufgerufen, wenn eine Aufgabe auf einen Fehler stößt. on_error ist optional. Implementiere es, um Fehler zu verwalten oder zu protokollieren.

Der Rückruf wird für jeden übermittelten Datensatz bzw. Batch einmal aufgerufen, sobald sein logischer Offset bestätigt wird oder die Verarbeitung fehlschlägt, und dient damit als fortlaufendes Signal für den Fortschritt der Datenerfassung über den gesamten Stream hinweg.

Deine Callback-Methoden werden auf den Hintergrundthreads des SDK ausgeführt, sodass ihr Aufruf deinen Producer nicht blockiert. Halte sie schnell und blockfrei. Was bei einem Fehler zu tun ist, liegt in der Verantwortung Ihres Kunden: protokollieren, alarmieren, erneut versuchen oder stoppen. Einige Fehler sind nicht behebbar, und wenn on_error meldet, dass der Stream endgültig fehlgeschlagen ist, musst du die Wiederherstellung auf einem neuen Stream durchführen. Siehe Erholungs- und Wiederholungsmuster.

Konfigurieren Sie einen Rückruf

Du hängst einen Callback an einen Stream an, indem du beim Erstellen des Streams eine Instanz deiner AckCallback-Unterklasse als ack_callback-Option in StreamConfigurationOptions übergibst. Der Rückruf gilt dann für jeden Datensatz, der in diesem Stream aufgenommen wird.

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

Wenn der Rückruf registriert ist, warten Sie nicht mehr inline. on_ack wird ausgelöst, wenn jeder Datensatz als dauerhaft bestätigt wird, und on_error wird ausgelöst, wenn ein Datensatz fehlschlägt.

Rückrufe im Vergleich zum Blockieren

Rückrufe und Blockierungsaufrufe lösen verschiedene Probleme, und Sie können sie zusammen verwenden:

  • Verwenden Sie einen Anerkennungsrückruf, um asynchron auf Dauerhaftigkeitsbestätigungen und Fehler zu reagieren, sobald sie auftreten, und dabei einen hohen Durchsatz aufrechtzuerhalten. Gut für Fortschrittsverfolgung, Kennzahlen und Fehlerprotokollierung.
  • Verwenden Sie wait_for_offset() oder flush(), wenn Ihr Code blockieren muss, bis ein bestimmter Datensatz oder alle ausstehenden Datensätze dauerhaft gespeichert wurden, bevor er fortfährt.