Meddelandeblockering och bekräftelse

Zerobus Ingest-SDK:erna erbjuder flera metoder för att mata in en post, som innebär en avvägning mellan genomströmning och hur mycket bekräftelse på varaktig lagring du får. Den här sidan förklarar varje metod och när man ska blockera hållbarheten. För att reagera asynkront på bekräftelser istället för blockerande, se bekräftelseåterkopplingar.

Exemplen på denna sida använder Python SDK. För exakt vilken timeout och konfigurationsalternativ varje metod accepterar (inklusive deras standardinställningar och enheter), se Zerobus SDK-arkivet. De andra språk-SDK:erna erbjuder motsvarande alternativ.

Vad är en offset?

Varje post du tar in tilldelas en offset: dess position i strömmen. Offset är hur du hänvisar till en specifik post när du vill bekräfta att den var varaktigt skriven. Zerobus Ingest ger en leveransgaranti av typen at-least-once, och att invänta en offset är det sätt en klient bekräftar den garantin för en specifik post.

Att bekräfta en offset innebär att posten är varaktigt lagrad, inte att den går att fråga efter i Delta-tabellen ännu. Zerobus Ingest skapar beständiga poster i tabellen strax därefter. För latenssiffror, se Latens.

Inmatningsmetoder

SDK:erna erbjuder två sätt att mata in en post. (Metodnamnen nedan kommer från Python SDK. Andra SDK:er exponerar motsvarande metoder.)

Metod Returns Använd den när
Offsetbaserad, ingest_record_offset() Posten är offset, efter att posten har köats i streamen. Rekommenderat standardval. Du vill köa poster i ordning och eventuellt bekräfta hållbarheten senare genom att vänta på en offset.
Framtidsbaserat, ingest_record() En RecordAcknowledgment du kan vänta på. Deprecated. Föredra offset-baserad för bättre prestanda.

Offset-baserad (rekommenderas)

ingest_record_offset() skickar in posten och returnerar dess offset när posten är köad i streamen. Anropet körs på den anropande tråden, så posterna köas i den ordning du anropar metoden, och den returnerade offseten gör att du senare kan verifiera persistens med wait_for_offset(). Detta är den rekommenderade standardinställningen för de flesta producenter, och det är metoden som används i Use Zerobus Ingest-exemplen .

Framtidsbaserad (föråldrad)

ingest_record() returnerar ett RecordAcknowledgment-objekt som du kan vänta på för att säkerställa beständighet. Den har föråldrats till förmån för offset-baserad metod, som presterar bättre. Använd det bara för befintlig kod som ännu inte har migrerats.

Post för post vs. batchinläsning

Varje inmatningsmetod har en batchvariant (till exempel ingest_records_offset()) som skickar en lista med poster i ett anrop. Batching är mer effektivt än enskilda anrop för bulkintagning.

För JSON och Protocol Buffers (protobuf) genomförs en batch atomärt: antingen accepteras varje post i batchen och lagras varaktigt, eller så avvisas hela batchen. Zerobus Ingest utför inte partiella uppladdningar eller partiell bekräftelse för dessa format, så din tabell innehåller aldrig en partiell batch. En batch som misslyckas med validering (till exempel en schema-mismatch) misslyckas snabbt, innan den rör vid tabellen, istället för att landa vissa poster och ta bort andra.

Eftersom en JSON- eller protobuf-batch skickas som ett enda meddelande, gäller den maximala meddelandestorleken på 10 MB både för en enskild post och en hel batch: alla poster i en batch tillsammans måste rymmas inom 10 MB. Anpassa storleken på dina batchar så att de håller sig under gränsen. Se Registerstorlek.

Arrow Flight-batcher utgör undantaget

Apache Arrow Flight-datainmatning följer inte modellen ovan med allt eller inget och ett enda meddelande. En Arrow-batch kan vara mycket större än en JSON- eller protobuf-batch, och Arrow Flight-vägen delar upp en stor batch i mindre transportmeddelanden som skickas och bekräftas individuellt istället för som en atomär enhet. Detta innebär följande:

  • Gränsen på 10 MB per meddelande som gäller för JSON- och protobuf-batchar gäller inte på samma sätt för en Arrow-batch. En stor Arrow-batch delas upp i transportmeddelanden i stället för att avvisas på grund av storleken.
  • Varaktighet bekräftas på transportmeddelandenivå, så en mycket stor logisk batch kan bli varaktig endast delvis om ett fel inträffar mitt i processen, i stället för att hela batchen genomförs enligt allt-eller-inget-principen.

ingest_batch() returnerar fortfarande en enda logisk offset för batchen du skickade in, och wait_for_offset() för den offseten slutförs först när varje transportmeddelande som utgör batchen har bekräftats. För hela Arrow Flight-modellen, batchstyrning och återställning av obekräftad data, se Use Arrow Flight with Zerobus Ingest.

När bör du blockera på ett meddelande?

Att blockera vid ett offset innebär att du byter genomströmning mot en starkare varaktighetsgaranti för varje post i din klientkod. Välj baserat på din arbetsbelastning:

  • Blockera inte: rätt standard för högvolymsströmning, där du bryr dig om uthållig genomströmning och kan bekräfta hållbarhet i aggregerad skala (till exempel vid strömstängning eller genom en bekräftelsecallback). De flesta producenter bör börja här.
  • Blockera vid en offset: tänk på detta när din applikation måste veta att en specifik post är varaktigt lagrad innan den utför någon annan åtgärd. Till exempel:
    • Du är på väg att ta bort eller bekräfta källan till datan (ett kömeddelande, en fil, en uppströms markör) och får inte förlora den om inmatningen misslyckas.
    • Du läser in data vid kontrollpunkter eller transaktionsgränser och behöver att varje kontrollpunkt är varaktigt lagrad innan du fortsätter.
    • Du gör skrivningar med låg volym men högt värde, där bekräftelse för varje post är viktigare än genomströmning.

Blockera inte på varje post i en höggenomströmningsslinga. Det tvingar din producent att vänta på ett tur och retur-anrop till servern för varje post och minskar genomströmningen kraftigt. Istället rekommenderar Azure Databricks att man matar in ett stort antal poster och sedan bekräftar varaktigheten en gång för hela mängden. Du har två sätt att göra det: vänta på den senaste offseten eller rensa strömmen. Blockering per enskild post bör reserveras för de specifika fallen ovan där en enskild post måste bekräftas innan nästa åtgärd.

Vänta på en offset

wait_for_offset() blockerar tills Zerobus Ingest bekräftar att posten vid den aktuella offseten har skrivits permanent, eller tills en timeout inträffar. Använd det för att bekräfta en specifik punkt i dataströmmen, oftast den sista posten i ett block. Mata in chunken, spara det slutliga offset-värdet som loopen returnerar, och vänta på just det offset-värdet i stället för att vänta efter varje post:

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

Spola bäcken

flush() blockerar tills alla poster som du hittills har matat in har skrivits till beständig lagring och returnerar sedan. Till skillnad från wait_for_offset(), spårar du inte en offset: flush väntar på allt som väntar på streamen. Den stänger inte bäcken, så du kan fortsätta äta efteråt.

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 vs. flush

Båda bekräftar hållbarhet för en bit. Välj baserat på vad du bekräftar:

  • Använd wait_for_offset(offset) när du vill bekräfta fram till en viss post, till exempel gränsen för en kontrollpunkt, medan andra poster fortfarande kan vara under bearbetning efter den.
  • Använd flush() när du vill bekräfta att alla väntande poster är hållbara innan du går vidare, till exempel i slutet av en batch, innan du flyttar en uppströms kursör eller innan du stänger ner. flush() styrs av en konfigurerbar flush-timeout.

close() spolar och stänger strömmen, så skivorna blir alltid hållbara vid en graciös avstängning. Anropa det alltid i ett finally-block.

Reagera asynkront på bekräftelser

Om du istället för att blockera vill reagera på hållbarhetsbekräftelser och fel när de kommer, medan din producent fortsätter på full fart, registrera en bekräftelseåterkoppling på streamen. Callbacks är en separat funktion från blockeringsanropen på denna sida. Se bekräftelseåteranrop.