Blockierung von Nachrichten und Bestätigung

Die Zerobus Ingest SDKs bieten mehrere Methoden zur Erfassung eines Datensatzes, die einen Kompromiss zwischen Durchsatz und dem Grad der Bestätigung der dauerhaften Speicherung eingehen, den Sie zurückerhalten. Diese Seite erklärt jede Methode und wann man bei der Dauerhaftigkeit blockieren sollte. Um asynchron auf Bestätigungen statt blockierend zu reagieren, siehe Bestätigungsrückrufe.

Die Beispiele auf dieser Seite verwenden das Python SDK. Für die genauen Timeout- und Konfigurationsoptionen, die jede Methode akzeptiert (einschließlich ihrer Standardwerte und Einheiten), siehe das Zerobus SDK-Repository. Die anderen Sprach-SDKs bieten entsprechende Optionen.

Was ist ein Offset?

Jedem Datensatz, den Sie in den Stream aufnehmen, wird ein Offset zugewiesen: seine Position im Stream. Der Offset ist die Art, sich auf einen bestimmten Eintrag zu beziehen, wenn Sie bestätigen möchten, dass er dauerhaft geschrieben wurde. Zerobus Ingest garantiert eine mindestens einmalige Zustellung und das Warten auf einen Offset ist der Weg, wie ein Client diese Garantie für einen bestimmten Datensatz bestätigt.

Das Bestätigen eines Offsets bedeutet, dass der Datensatz dauerhaft gespeichert ist, nicht dass er bereits in der Delta-Tabelle abgefragt werden kann. Zerobus Ingest materialisiert kurz darauf dauerhafte Datensätze in der Tabelle. Für Latenzzahlen siehe Latenz.

Erfassungsmethoden

Die SDKs bieten zwei Möglichkeiten, einen Datensatz einzulesen. (Die untenstehenden Methodennamen stammen aus dem Python SDK. Andere SDKs bieten vergleichbare Methoden an.)

Methode Rückgaben Verwenden Sie ihn, wenn
Offset-basiert, ingest_record_offset() Der Offset des Datensatzes, nachdem der Datensatz in den Stream eingereiht wurde. Empfohlene Standardeinstellung. Sie sollten Datensätze der Reihe nach einreihen und optional die Dauerhaftigkeit später bestätigen, indem Sie auf einen Offset warten.
Zukunftsbasiert, ingest_record() Ein RecordAcknowledgment, auf das Sie warten können. Deprecated. Bevorzugen Sie die offsetbasierte Methode für eine bessere Leistung.

Offsetbasiert (empfohlen)

ingest_record_offset() übermittelt den Datensatz und gibt dessen Offset zurück, sobald der Datensatz in die Warteschlange des Streams eingereiht wurde. Der Aufruf läuft auf Ihrem aufrufenden Thread, sodass Datensätze in der Reihenfolge des Aufrufs der Methode in die Warteschlange eingereiht werden, und der zurückgegebene Offset ermöglicht es Ihnen, die Dauerhaftigkeit später mit wait_for_offset() zu bestätigen. Dies ist der empfohlene Standard für die meisten Hersteller und die Methode, die in den Use Zerobus Ingest-Beispielen verwendet wird.

Zukunftsorientiert (veraltet)

ingest_record() Gibt ein RecordAcknowledgment Objekt zurück, auf das du für die Haltbarkeit warten kannst. Sie wird zugunsten der offsetbasierten Methode eingestellt, die besser abschneidet. Verwende es nur für bestehenden Code, der noch nicht migriert wurde.

Datensatzweise Erfassung im Vergleich zu Batch-Erfassung

Jede Eingabemethode hat eine Batch-Variante (zum Beispiel ingest_records_offset()), die eine Liste von Datensätzen in einem Aufruf einreicht. Stapelverarbeitung ist effizienter als einzelne Aufrufe für die Massenerfassung.

Für JSON und Protocol Buffers (protobuf) wird ein Batch atomar festgeschrieben: Entweder wird jeder Datensatz im Batch akzeptiert und dauerhaft gespeichert oder der gesamte Batch wird abgelehnt. Zerobus Ingest führt keine teilweisen Uploads oder teilweisen Bestätigungen für diese Formate durch, sodass Ihre Tabelle nie einen teilweisen Batch enthält. Ein Batch, der bei der Validierung fehlschlägt (zum Beispiel ein Schema-Missmatch), scheitert schnell, bevor er die Tabelle berührt, anstatt einige Datensätze zu landen und andere zu verwerfen.

Da ein JSON- oder Protobuf-Batch als einzelne Nachricht gesendet wird, gilt die maximale Nachrichtengröße von 10 MB sowohl für einen einzelnen Datensatz als auch für einen ganzen Batch: Alle Datensätze in einem Batch zusammen müssen innerhalb von 10 MB passen. Dimensioniere deine Chargen so, dass sie unter diesem Limit bleiben. Siehe Datensatzgröße.

Arrow-Flight-Batches sind die Ausnahme

Die Datenaufnahme über Apache Arrow Flight folgt nicht dem oben beschriebenen Alles-oder-Nichts-Modell mit einer einzelnen Nachricht. Ein Arrow-Batch kann viel größer sein als ein JSON- oder Protubuf-Batch, und der Arrow Flight-Pfad teilt einen großen Batch in kleinere Transport-Nachrichten auf, die einzeln gesendet und bestätigt werden, anstatt als eine atomare Einheit. Dies führt zu folgendem Ergebnis:

  • Das 10-MB-Limit pro Nachricht, das für JSON- und Protobuf-Batches gilt, gilt nicht in gleicher Weise für einen Arrow-Batch. Ein großer Arrow-Batch wird in Transportnachrichten unterteilt, anstatt aufgrund seiner Größe abgelehnt zu werden.
  • Die Dauerhaftigkeit wird bei der Transport-Nachrichten-Granularität bestätigt, sodass ein sehr großer logischer Batch teilweise dauerhaft sein kann, wenn ein Fehler mitten im Verlauf auftritt, anstatt Alles-oder-Nichts zu committen.

ingest_batch() gibt weiterhin einen einzelnen logischen Offset für den von Ihnen eingereichten Batch zurück und wait_for_offset() wird bei diesem Offset erst abgeschlossen, nachdem jede Transportnachricht, die den Batch ausmacht, anerkannt wurde. Für das vollständige Modell von Arrow Flight, Hinweise zum Batching und die Wiederherstellung nicht bestätigter Daten, siehe Use Arrow Flight with Zerobus Ingest.

Wann sollte man eine Nachricht blockieren?

Das Blockieren bei einem Offset tauscht Durchsatz gegen eine stärkere Dauerhaftigkeitsgarantie pro Datensatz in Ihrem Client-Code ein. Wählen Sie basierend auf Ihrer Arbeitsbelastung:

  • Nicht blockieren: der richtige Standard für hochvolumiges Streaming, bei dem dir ein nachhaltiger Durchsatz wichtig ist und du die Dauerhaftigkeit aggregiert bestätigen kannst (zum Beispiel beim Stream-Schließen oder durch einen Bestätigungsrückruf). Die meisten Produzenten sollten hier anfangen.
  • Bei einem Offset blockieren: Beachten Sie dies, wenn Ihre Anwendung wissen muss, dass ein bestimmter Datensatz dauerhaft ist, bevor sie eine weitere Aktion ausführt. Beispiel:
    • Sie stehen kurz davor, die Quelle der Daten (eine Warteschlangennachricht, eine Datei, einen Upstream-Cursor) zu löschen oder anzuerkennen, und dürfen sie nicht verlieren, falls die Aufnahme fehlschlägt.
    • Sie erfassen Daten an Prüfpunkten oder Transaktionsgrenzen und müssen sicherstellen, dass jeder Prüfpunkt dauerhaft ist, bevor Sie fortfahren.
    • Sie führen Schreibvorgänge in geringem Umfang, aber mit hohem Wert durch, bei denen die Bestätigung jedes einzelnen Datensatzes wichtiger ist als der Durchsatz.

Blockieren Sie nicht bei jedem Datensatz in einer Schleife mit hohem Durchsatz. Das serialisiert Ihren Produzenten auf eine Roundtrip zum Server für jeden Datensatz und reduziert den Durchsatz stark. Stattdessen empfiehlt Azure Databricks, einen großen Block an Datensätzen aufzunehmen und dann die dauerhafte Speicherung einmal für den gesamten Block zu bestätigen. Du hast zwei Möglichkeiten, das zu tun: auf den neuesten Offset warten oder den Stream auslöschen. Das Blockieren pro einzelner Datensatz sollte für die oben genannten speziellen Fälle reserviert sein, in denen ein einzelner Datensatz vor der nächsten Aktion bestätigt werden muss.

Auf einen Offset warten

wait_for_offset() blockiert, bis Zerobus Ingest bestätigt, dass der Datensatz an diesem Offset dauerhaft geschrieben wurde, oder bis eine Zeitüberschreitung eintritt. Verwenden Sie dies, um einen bestimmten Punkt in einem Stream zu bestätigen, meist den letzten Datensatz eines Segments. Erfassen Sie das Segment, behalten Sie den letzten Offset, den die Schleife zurückgibt, und warten Sie auf diesen einen Offset, anstatt nach jedem Datensatz zu warten:

from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

table_properties = TableProperties("main.default.air_quality")
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

try:
    last_offset = 0
    for row in records:
        last_offset = stream.ingest_record_offset(row)

    # Block until everything up to the last record of the chunk is durable
    stream.wait_for_offset(last_offset)
    print("Chunk durably written.")
finally:
    stream.close()

Leeren Sie den Stream

flush() blockiert, bis alle Datensätze, die Sie bislang erfasst haben, dauerhaft geschrieben sind, und gibt dann die Kontrolle zurück. Im Gegensatz zu wait_for_offset() wird ein Offset nicht verfolgt: Flush wartet auf alles, was im Stream aussteht. Es schließt den Stream nicht, sodass Sie danach weiter erfassen können.

try:
    for row in records:
        stream.ingest_record_offset(row)

    # Block until every pending record is durable
    stream.flush()
    print("All ingested records durably written.")
finally:
    stream.close()

wait_for_offset im Vergleich zu Flush

Beide bestätigen die Dauerhaftigkeit für ein Segment. Wählen Sie basierend auf dem, was Sie bestätigen möchten:

  • Verwenden Sie wait_for_offset(offset), wenn Sie bis zu einem bestimmten Datensatz bestätigen möchten, zum Beispiel eine Prüfpunktgrenze, während andere Datensätze noch dahinter in Flight sind.
  • Verwenden Sie flush(), wenn Sie sicherstellen möchten, dass alle ausstehenden Datensätze dauerhaft sind, bevor Sie fortfahren, zum Beispiel am Ende eines Batches, bevor Sie einen Upstream-Cursor weiterschieben oder bevor Sie das System herunterfahren. flush() wird durch ein konfigurierbares Flush-Timeout bestimmt.

close() spült und schließt den Stream, sodass Datensätze bei ordnungsgemäßem Herunterfahren immer dauerhaft sind. Rufen Sie es immer innerhalb eines finally-Blocks auf.

Asynchron auf Bestätigungen reagieren

Wenn sie statt zu blockieren auf Dauerhaftigkeitsbestätigungen und Fehler reagieren möchten, während Ihr Producer mit voller Geschwindigkeit weiterpusht, registrieren Sie eine Anerkennungsrückmeldung im Stream. Callbacks sind eine separate Funktion und unterscheiden sich von den blockierenden Aufrufen auf dieser Seite. Siehe Rückrufe zur Bestätigung.