Asynchrone communicatie

De communicatie op een stroom is asynchroon en bidirectioneel. Je client pusht records continu zonder te wachten tot ze allemaal bevestigd zijn, en de server stuurt bevestigingen terug via dezelfde verbinding zodra records duurzaam worden. Deze ontkoppeling maakt het mogelijk dat één client een hoge doorvoersnelheid kan aanhouden: hij blijft gegevens verzenden terwijl ontvangstbevestigingen op de achtergrond binnenkomen.

Asynchrone client- en servercommunicatie op een Zerobus Ingest-stream: de client pusht continu records terwijl de server committed offsets terugstuurt over dezelfde bidirectionele verbinding naarmate records duurzaam worden

Offsets en de bevestigingscyclus

Elke indiening op een stream, of het nu een enkel record of een batch is, krijgt een logische offset toegewezen die zijn positie in die stream aangeeft. In plaats van elke aanlevering afzonderlijk te bevestigen, rapporteert de server de cumulatieve duurzaamheidsstatus via de hoogste vastgelegde offset die hij tot dusver duurzaam heeft opgeslagen. Omdat offsets geordend zijn, bevestigt één bevestiging die inzending en alle eerdere inzendingen.

Dit is de bevestigingslus, en die houdt de verbinding zowel snel als betrouwbaar:

  1. De client stuurt records en bewaart ze in een lokale in-flight-buffer.
  2. De server bewaart records duurzaam en stuurt periodiek de hoogste gecommitteerde offset terug.
  3. Na ontvangst van die offset verwijdert de cliënt veilig elk gebufferd record tot aan de cliënt, omdat die records nu duurzaam zijn.

Wanneer je een Zerobus Ingest SDK gebruikt, voert de SDK deze lus voor je uit. Het houdt offsets bij, onderhoudt de buffer tijdens de vlucht en verwerkt bevestigingen op de achtergrond terwijl je producer blijft pushen. Je implementeert de lus niet zelf. Wat je optioneel kunt bepalen is hoe je de duurzaamheid waarneemt:

  • Blijf sturen; de SDK verwerkt bevestigingen zodra ze binnenkomen.
  • Blokkeer een offset alleen wanneer je aanvraag moet wachten tot een specifiek record duurzaam is. Zie hieronder.
  • Registreer een bevestigingscallback om asynchroon te reageren op bevestigingen en fouten, zonder blokkering. Zie Acknowledgment callbacks.

Je zou de offset-tracking en buffering loop alleen zelf implementeren als je een custom client bouwt die geen SDK gebruikt.

De in-flight buffer wordt begrensd door een configureerbare limiet voor in-flight records. Inname verloopt asynchroon totdat de buffer vol raakt; op dat moment blokkeren inname-aanroepen totdat bevestigingen binnenkomen en er weer ruimte vrijkomt. Stel de limiet voor je werklast af en let op dat gebufferde records het geheugen van de client verbruiken terwijl ze in flight zijn. Voor de optie en de standaardinstelling, zie de Zerobus SDK-repository.

Als de verbinding wordt onderbroken, zijn records die nog in de in-flight buffer zitten (die na de laatste committed offset) niet als duurzaam bevestigd, dus kunnen ze opnieuw worden afgespeeld. Zie Herstel- en herpogingspatronen.

De bevestiging bevestigt duurzaamheid, niet de aanvraagbaarheid. Een gecommitteerde compensatie betekent dat die records duurzaam worden bewaard en niet verloren zullen gaan. Zerobus Ingest materialiseert duurzame records kort daarna als een aparte stap in de Delta-tabel, waarna de data binnen ongeveer 5 seconden opvraagbaar wordt. Voor meer informatie over latency, zie Latency.

Wachten op een record versus maximale doorvoer maximaliseren

Je wacht op de offset van een record wanneer je applicatie verdere uitvoering moet blokkeren totdat dat specifieke record als duurzaam bekend is, bijvoorbeeld voordat het werk wordt bevestigd aan een upstream-systeem. Wachten draait om synchronisatie op applicatieniveau, niet om duurzaamheid. Een record wordt duurzaam via de bevestigingslus, ongeacht of je erop blokkeert of niet.

Blokkeren heeft een doorvoerkosten:

  • Wachten na elk record maakt van de gegevensinname een feitelijk synchrone workflow. Blokkeren op elk bericht voordat het volgende wordt verzonden, voorkomt dat de client de volledige doorvoer van Zerobus Ingest bereikt.
  • Gegevensinname met hoge doorvoersnelheid is continu en asynchroon. De client blijft records verzenden terwijl bevestigingen binnenkomen voor groepen eerdere records, in plaats van bij elk record te pauzeren. Wacht op een specifieke offset alleen bij de checkpoints waar je applicatie die garantie echt nodig heeft, of gebruik een bevestigingscallback om de voortgang te volgen zonder te blokkeren.

Voor de invoermethoden, wanneer te blokkeren bij een offset en hoe bevestigingscallbacks werken, zie Bericht blokkeren en bevestigen.

Volgorde in een gegevensstroom

Bevestigingen en offsets zijn per stream: de volgorde is gegarandeerd binnen één enkele stream, niet over alle streams heen. Voor hoe per-stream ordering werkt en hoe je eromheen ontwerpt, zie Ordering guarantees.