Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
De Zerobus Ingest SDK’s bieden verschillende methoden om een record op te nemen, waarbij er een afweging wordt gemaakt tussen de doorvoersnelheid en de mate van bevestiging van duurzaamheid die u ontvangt. Deze pagina legt elke methode uit en wanneer je moet blokkeren op duurzaamheid. Om asynchroon op bevestigingen te reageren in plaats van te blokkeren, zie Acknowledgment callbacks.
De voorbeelden op deze pagina gebruiken de Python SDK. Voor de exacte timeout en configuratieopties die elke methode accepteert (inclusief hun standaardwaarden en eenheden), zie de Zerobus SDK-repository. De SDK's van de andere taal bieden equivalente opties.
Wat is een offset?
Elke registratie die je invoert krijgt een offset toegewezen: de positie in de stroom. De offset is hoe je naar een specifiek document verwijst wanneer je wilt bevestigen dat het duurzaam is geschreven. Zerobus Ingest biedt minstens één leveringsgarantie, en wachten op een offset is hoe een client die garantie voor een bepaald record bevestigt.
Het bevestigen van een offset betekent dat het record duurzaam is, niet dat het al in de Delta-tabel kan worden opgevraagd. Zerobus Ingest schrijft persistente records kort daarna naar de tabel. Voor latentiecijfers, zie Latency.
Opnamemethoden
De SDK's bieden twee manieren om een record te importeren. (Methodennamen hieronder komen uit de Python SDK. Andere SDK's bieden vergelijkbare methoden.)
| Method | Returns | Gebruik deze wanneer |
|---|---|---|
Op basis van offset, ingest_record_offset() |
De offset van het record, nadat het record in de stream is geplaatst. | Aanbevolen standaard. Je wilt records in volgorde in de wachtrij zetten en eventueel later de duurzaamheid bevestigen door te wachten op een offset. |
Toekomstgebonden, ingest_record() |
Een RecordAcknowledgment waarop je kunt wachten. |
Deprecated. Geef de voorkeur aan offset-gebaseerd voor betere prestaties. |
Op offset gebaseerd (aanbevolen)
ingest_record_offset() verzendt het record en retourneert de offset zodra het record in de wachtrij van de stream is geplaatst. De aanroep wordt uitgevoerd op je aanroepthread, dus records worden in de wachtrij gezet in de volgorde waarin je de methode aanroept, en met de geretourneerde offset kun je later de duurzaamheid bevestigen met wait_for_offset(). Dit is de aanbevolen standaard voor de meeste producenten, en het is de methode die wordt gebruikt in de Use Zerobus Ingest-voorbeelden .
Op futures gebaseerd (verouderd)
ingest_record() Geeft een RecordAcknowledgment object terug waar je op kunt wachten voor duurzaamheid. Deze is afgeschaft ten gunste van de offset-gebaseerde methode, die beter presteert. Gebruik het alleen voor bestaande code die nog niet is gemigreerd.
Record-voor-record vs. batch-ingestie
Elke ingestiemethode heeft een batchvariant (bijvoorbeeld ingest_records_offset()) die een lijst met records in één aanroep indient. Batching is efficiënter dan individuele oproepen voor bulkopname.
Voor JSON en Protocol Buffers (protobuf) geldt dat een batch atomair wordt vastgelegd: ofwel wordt elk record in de batch geaccepteerd en duurzaam opgeslagen, ofwel wordt de hele batch afgewezen. Zerobus Ingest voert geen gedeeltelijke uploads of gedeeltelijke bevestiging uit voor deze formaten, dus je tabel bevat nooit een gedeeltelijke batch. Een batch die de validatie niet doorstaat (bijvoorbeeld door een schema-mismatch) mislukt direct, voordat die in de tabel terechtkomt, in plaats van sommige records weg te schrijven en andere te verwerpen.
Omdat een JSON- of protobuf-batch als één bericht wordt verzonden, geldt de maximale berichtgrootte van 10 MB voor zowel één enkel record als een hele batch: alle records in een batch samen moeten binnen 10 MB passen. Pas je batches zo aan dat het binnen die limiet blijft. Zie grootte van de record.
Arrow Flight-batches zijn de uitzondering
Apache Arrow Flight-gegevensinname volgt niet het hierboven beschreven alles-of-nietsmodel van één enkel bericht. Een Arrow-batch kan veel groter zijn dan een JSON- of protobuf-batch, en het Arrow Flight-pad splitst een grote batch op in kleinere transportberichten die individueel worden verzonden en bevestigd in plaats van als één atomaire eenheid. Als gevolg hiervan:
- De limiet van 10 MB per bericht die geldt voor JSON- en protobuf-batches geldt niet op dezelfde manier voor een Arrow-batch. Een grote Arrow-batch wordt verdeeld in transportberichten in plaats van afgewezen te worden vanwege de grootte.
- Duurzame opslag wordt bevestigd op het niveau van afzonderlijke transportberichten, waardoor een zeer grote logische batch slechts gedeeltelijk duurzaam kan zijn opgeslagen als er tijdens de verwerking een storing optreedt, in plaats van in zijn geheel of helemaal niet te worden vastgelegd.
ingest_batch() Geeft nog steeds één logische offset terug voor de batch die je hebt ingediend, en wait_for_offset() op die offset wordt deze pas voltooid nadat elk transportbericht dat de batch vormt is bevestigd. Voor het volledige Arrow Flight-model, batchbegeleiding en het herstel van niet-erkende data, zie Use Arrow Flight with Zerobus Ingest.
Wanneer moet je een bericht blokkeren?
Blokkeren op een offset gaat ten koste van de doorvoer ten gunste van een sterkere garantie op persistentie per record in je clientcode. Kies op basis van je werkdruk:
- Niet blokkeren: de juiste standaard voor streaming met hoog volume, waarbij je geeft om een blijvende doorvoer en duurzaamheid in totaal kunt bevestigen (bijvoorbeeld bij het sluiten van de stream of via een bevestigingscallback). De meeste producenten zouden hier moeten beginnen.
-
Blokkeren bij een offset: overweeg dit wanneer je aanvraag moet weten dat een specifiek record duurzaam is voordat er een nieuwe actie wordt ondernomen. Bijvoorbeeld:
- Je staat op het punt de bron van de data te verwijderen of te bevestigen (een wachtrijbericht, een bestand, een cursor stroomopwaarts) en mag deze niet verliezen als de invoer mislukt.
- Je neemt checkpoints of transactionele grenzen binnen en je wilt dat elk checkpoint duurzaam is voordat je verder gaat.
- Je doet schrijfopdrachten met een laag volume en hoge waarde waarbij bevestiging per record belangrijker is dan doorvoer.
Blokkeer niet op elk record in een high-throughput lus. Daardoor wordt je producer voor elk record gedwongen tot een heen-en-weerverzoek naar de server, wat de doorvoersnelheid sterk verlaagt. In plaats daarvan raadt Azure Databricks aan om een groot deel van de records te importeren en vervolgens de duurzaamheid eenmaal voor het hele stuk te bevestigen. Je hebt twee manieren om dat te doen: wachten op de laatste offset, of de stream spoelen. Blokkeren per individueel record moet worden voorbehouden aan de specifieke gevallen hierboven waarin één record bevestigd moet worden vóór de volgende actie.
Wacht op een offset
wait_for_offset() blokkeert totdat Zerobus Ingest bevestigt dat het record bij die offset duurzaam is geschreven, of totdat het timeout is. Gebruik het om een specifiek punt in de stream te bevestigen, meestal het laatste record van een stuk. Voer de chunk in, houd de laatste offset die de loop teruggeeft en wacht op die ene offset in plaats van na elke record te wachten:
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()
Spoel de stream leeg
flush() blokkeert totdat alle records die u tot nu toe hebt opgenomen duurzaam zijn opgeslagen, en keert vervolgens terug. In tegenstelling tot wait_for_offset(), houd je geen offset bij: flush wacht op alles wat in de stream hangt. Het sluit de stroom niet, dus je kunt daarna blijven eten.
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 versus flush
Beide bevestigen de duurzaamheid voor geruime tijd. Kies op basis van wat je wilt bevestigen:
- Gebruik
wait_for_offset(offset)wanneer je tot en met een specifiek record wilt bevestigen, bijvoorbeeld een checkpointgrens, terwijl andere records daarachter mogelijk nog onderweg zijn. - Gebruik
flush()wanneer je zeker wilt weten dat alle openstaande records duurzaam zijn opgeslagen voordat je verdergaat, bijvoorbeeld aan het einde van een batch, voordat je een upstream-cursor vooruitzet, of voordat je afsluit.flush()wordt bepaald door een configureerbare flush-time-out.
close() spoelt en sluit de stream, zodat records altijd duurzaam worden bij een elegante afsluiting. Roep het altijd aan binnen een finally-blok.
Reageer asynchroon op bevestigingen
Als je, in plaats van te blokkeren, wilt reageren op duurzaamheidsbevestigingen en -fouten zodra die binnenkomen, terwijl je producer op volle snelheid berichten blijft versturen, registreer dan een acknowledgment-callback op de stream. Callbacks zijn een aparte functionaliteit ten opzichte van de blokkeringsoproepen op deze pagina. Zie Acknowledgment callbacks.
Related
- Bevestigingscallbacks: Reageer asynchroon op bevestigingen en fouten.
- Gebruik Zerobus Ingest: Schrijf een clienttoepassing.
- Streams: Gegevensstromen en offsets.
- Zerobus Ingest foutafhandeling: Foutafhandeling.