7.6

View in English

7.6 Realtime- en streamingdata

Overzicht en motivatie

Het meeste van wat je over datapijplijnen weet gaat ervan uit dat de data stilstaat. Je verzamelt een dag records, draait ‘s nachts een job en leest ‘s ochtends de resultaten. Realtime- en streamingdata keert die aanname om. In plaats van een afgeronde stapel data te verwerken, verwerk je een eindeloze stroom gebeurtenissen terwijl ze aankomen, en produceer je continu antwoorden. Dit is het verschil tussen batchverwerking, die werkt op een begrensde, complete dataset, en streamverwerking, die werkt op een onbegrensde, nooit afgeronde stroom.

Voor grote teams verschijnt streaming zodra latentie voor het bedrijf telt. Een fraudebeslissing die een uur te laat komt is waardeloos. Een personalisatiesignaal dat morgen landt personaliseert niets. Een operationeel dashboard dat een dienst achterloopt op de werkelijkheid misleidt de mensen die ernaar kijken. Hoofdstuk 7.2 (data-engineering) betoogt dat je standaard batch moet kiezen en alleen naar streaming moet grijpen waar latentie werkelijk loont, en dit hoofdstuk brengt je de rest van de weg: wanneer realtime zijn kosten verdient, en hoe je het bouwt zonder je operationele budget in brand te steken. Streaming zit dicht bij de gebeurtenisgedreven berichtenpatronen in hoofdstuk 3.12 (gebeurtenisgedreven architectuur en berichtenverkeer), de opslagkeuzes in hoofdstuk 3.4 (dataarchitectuur en opslag) en de telemetriepraktijken in hoofdstuk 9.2 (observeerbaarheid en telemetrie).

Omgevingen van onderneming en overheid verhogen de inzet. Een bank scoort transacties op fraude in de tijd die een kaartlezer nodig heeft om te knipperen. Een vervoersinstantie volgt voertuigen en voorspelt aankomsten voor miljoenen reizigers. Een uitkeringsinstantie let op anomalieën in claims terwijl ze een controleerbaar register van elke beslissing bijhoudt. In al deze gevallen komt de waarde uit handelen op data terwijl ze nog vers is, en het risico uit handelen op data die fout, onvolledig of achteraf niet te reconstrueren is. Dit hoofdstuk is uitgesproken over beide.

Kernprincipes

  • Grijp alleen naar streaming wanneer latentie een duidelijke bedrijfswaarde heeft. Batch is goedkoper en eenvoudiger.
  • Onderscheid begrensde (eindige) data van onbegrensde (eindeloze) data en ontwerp dienovereenkomstig.
  • Behandel gebeurtenistijd, niet aankomsttijd, als bron van waarheid en plan voor late en niet-geordende data.
  • Vensters en watermarks zijn hoe je eindige antwoorden uit oneindige stromen krijgt.
  • Geef de voorkeur aan effectief-eenmalige resultaten via idempotente sinks boven fragiele exactly-oncebeloftes.
  • Toestandsgebonden verwerking heeft checkpointing nodig zodat ze kan herstellen zonder te verliezen of dubbel te tellen.
  • Ontwerp vanaf dag één voor tegendruk en herverwerking, niet als bijgedachte.
  • Houd streaminglogica observeerbaar en controleerbaar. Een stille stroom is erger dan een gefaalde batch.

Aanbevelingen

Rechtvaardig realtime voordat je het bouwt

De belangrijkste streamingbeslissing is of je überhaupt streamt. Realtime verdubbelt ruwweg je operationele complexiteit en kosten, omdat je een job die draait en stopt ruilt tegen een systeem dat elke seconde gezond moet blijven. Noem voordat je je vastlegt de beslissing die verse data mogelijk maakt en de kosten van die beslissing die te laat aankomt. Fraudescoring, operationele alarmering en live personalisatie halen de lat meestal. Een dashboard dat een mens twee keer per dag bekijkt bijna nooit, hoe bevredigend “realtime” ook klinkt in een planningsvergadering. Schrijf de latentie-eis op als getal, in seconden of minuten, en toets haar aan de werkelijkheid. Veel van wat mensen realtime noemen wordt goed bediend door microbatches die om de paar minuten draaien tegen een fractie van de kosten.

Ontwerp rond gebeurtenistijd, niet verwerkingstijd

Het moeilijkste idee in streaming is dat gebeurtenissen op het ene moment plaatsvinden en op een ander worden verwerkt. Gebeurtenistijd is wanneer het ding werkelijk gebeurde, bijvoorbeeld wanneer een reiziger een kaart tikte. Verwerkingstijd is wanneer je systeem eraan toekwam. Deze lopen voortdurend uiteen: een telefoon verliest signaal in een tunnel en uploadt drie minuten tikken ineens, een netwerkhapering herordent berichten, een partitie loopt achter. Als je op verwerkingstijd rekent, wiebelen je getallen met je infrastructuur in plaats van de wereld te weerspiegelen. Dit probleem van laat aankomende en niet-geordende gebeurtenissen is het hart van de discipline, en het sluit direct aan op de gebeurtenismodellering in gebeurtenisgedreven architectuur. Stempel elke gebeurtenis bij de bron met haar gebeurtenistijd, draag die tijdstempel door de hele pijplijn en bereken je resultaten ertegen.

Gebruik vensters en watermarks voor eindige antwoorden

Een onbegrensde stroom eindigt nooit, dus “tel de gebeurtenissen” heeft geen antwoord tot je haar begrenst. Vensters doen die begrenzing. Tumbling-vensters hakken tijd in vaste, niet-overlappende emmers, bijvoorbeeld elke minuut. Sliding-vensters overlappen, dus een venster van vijf minuten dat elke minuut opschuift geeft je een vloeiend voortschrijdend cijfer. Sessievensters groeperen bursts van activiteit gescheiden door gaten van inactiviteit, wat goed past bij gebruikerssessies. Zodra je vensters hebt, moet je beslissen wanneer een venster klaar is, want late data kan nog aankomen. Een watermark is de schatting van het systeem dat het waarschijnlijk alle gebeurtenissen tot een gegeven gebeurtenistijd heeft gezien. Wanneer de watermark het einde van een venster passeert, geef je het resultaat af. Stem af hoe lang je wacht: houd vensters langer open en je verdraagt meer vertraging ten koste van latentie en geheugen, sluit ze sneller en je riskeert achterblijvers te laten vallen. Besluit expliciet wat er gebeurt met data die aankomt nadat een venster sluit, of je haar laat vallen, logt of een correctie afgeeft.

Maak sinks idempotent en geef de voorkeur aan effectief-eenmalig

Leveringsgaranties klinken eenvoudig en zijn dat niet. At-least-oncelevering betekent dat elke gebeurtenis wordt verwerkt, maar sommige na een herpoging meer dan eens, zodat tellingen kunnen opzwellen. Exactly-once klinkt ideaal maar is duur en, letterlijk genomen over willekeurige externe systemen, vaak onmogelijk. Het praktische doel is effectief-eenmalig: het waarneembare resultaat is alsof elke gebeurtenis eenmaal werd verwerkt, ook al herprobeerde de machinerie eronder. Je komt daar door je idempotente sinks veilig te maken om herhaaldelijk naar te schrijven, met deterministische sleutels en upserts zodat een herspeelde gebeurtenis overschrijft in plaats van dupliceert. Combineer at-least-oncelevering met idempotente schrijfacties en je krijgt correcte resultaten zonder overal te betalen voor zware transactiecoördinatie. Reserveer echte exactly-oncemachinerie voor de smalle plekken die het werkelijk nodig hebben.

Checkpoint toestandsgebonden verwerking zodat ze kan herstellen

Veel nuttige streamingberekeningen zijn toestandsgebonden: lopende tellingen, joins over stromen, ontdubbelen, fraudemodellen die recent gedrag onthouden. Die toestand leeft in het geheugen en zou verdwijnen wanneer een proces herstart. Checkpointing maakt periodiek een momentopname van de toestand en de streampositie samen, zodat het systeem na een crash vanaf een consistent punt hervat in plaats van alles opnieuw af te spelen of zijn geheugen te verliezen. Dimensioneer je toestand bewust, want onbegrensde toestand is een gangbare manier om een streamingjob in productie zonder geheugen te laten raken. Gebruik verloop en time-to-live op toestand die je niet langer nodig hebt en bewaak toestandsgrootte als eersterangs statistiek. Hersteltijd na een falen is een echte service-levelzorg, dus test haar voordat je gebruikers dat doen.

Stream vanuit operationele databases met change data capture

Je wilt vaak reageren op wijzigingen in een database die nooit is ontworpen om gebeurtenissen uit te zenden. Change data capture (CDC) lost dit op door het transactielogboek van de database te lezen en elke insert, update en delete om te zetten in een stroom wijzigingsgebeurtenissen. Dit is veel beter dan de tabel op een timer pollen, wat traag is, tussentoestanden mist en de bron bestookt. CDC laat je een zoekindex, een cache, een analyticsopslag of een downstreamservice continu synchroon houden met een systeem van registratie, en doet dat zonder ingrijpende wijzigingen in de applicatie. Behandel de wijzigingsstroom als eersterangs dataproduct: versioneer haar schema, documenteer haar betekenis en bewaak haar vertraging, want alles stroomafwaarts erft die vertraging.

Geef de voorkeur aan een streaming-first-architectuur boven twee codebases onderhouden

De klassieke Lambda-architectuur draait een batchlaag voor nauwkeurige, complete geschiedenis naast een snelheidslaag voor verse, benaderende resultaten, en voegt ze dan samen. Het werkt, maar het laat je dezelfde bedrijfslogica tweemaal schrijven en onderhouden, in twee systemen, en de verschillen eeuwig afstemmen. De Kappa-architectuur klapt dit samen: houd een duurzaam, herspeelbaar logboek van gebeurtenissen en draai alle verwerking als streamverwerking, geschiedenis herverwerkend door het logboek opnieuw af te spelen wanneer logica verandert. De sector is naar deze streaming-firstvorm afgedreven omdat één codebasis dramatisch goedkoper is om te onderhouden en over te redeneren. Als je je batchbehoeften kunt uitdrukken als herspelingen over een bewaard gebeurtenislogboek, vermijd je de belasting van twee codebases volledig. Gebruik op logboeken gebaseerde brokers die geschiedenis bewaren zodat herverwerking een kwestie is van terugspoelen, niet herbouwen.

Stel streams beschikbaar als SQL, gematerialiseerde views en realtime-OLAP

Niet iedereen die streaming nodig heeft zou laag-niveaustreamverwerkingscode moeten schrijven. Streaming-SQL laat analisten en engineers vensters, joins en aggregaties uitdrukken in een taal die ze al kennen, en houdt de resultaten continu actueel als gematerialiseerde views. Voor analytische queries met lage latentie over verse data neemt een realtime-online analytical processing (OLAP)-opslag de stroom in en beantwoordt slice-and-dicequeries in milliseconden, wat een werkelijk live operationeel dashboard aandrijft. Koppel deze aan de productanalyticspraktijken in hoofdstuk 7.4 (productanalytics en experimenteren) wanneer het doel snelle feedback op functies en experimenten is. Kies deze hogere-niveautools waar ze passen en bewaar met de hand geschreven streamprocessors voor logica die ze niet kunnen uitdrukken.

Plan vanaf het begin voor tegendruk en herverwerking

Een stroom kan sneller aankomen dan je kunt verwerken. Tegendruk (backpressure) is het mechanisme waarmee een trage afnemer stroomopwaarts kan signaleren te vertragen in plaats van om te vallen of stilletjes data te laten vallen. Zorg dat elke fase in je pijplijn het eert en bewaak afnemersvertraging (lag) als kopstatistiek, want groeiende lag is de vroegste waarschuwing dat je de race verliest. Herverwerking is het andere vermogen dat mensen wensen te hebben ingebouwd. Wanneer je een bug vindt of een regel wijzigt, wil je geschiedenis door de gecorrigeerde logica herspelen. Dat kan alleen als je gebeurtenislogboek genoeg geschiedenis bewaart en je sinks idempotent genoeg zijn om de herspeling te absorberen. Ontwerp beide vanaf dag één. Ze achteraf aanbrengen onder incidentdruk is miserabel.

Afwegingen: voor- en nadelen

KeuzeVoordelenNadelenPast het best bij
BatchEenvoudig, goedkoop, makkelijk te testen en backfillenHoge latentie, verouderd tussen runsRapportage, de meeste analytics
Microbatch (minuten)Bijna realtime, veel eenvoudiger dan streamingNiet echt direct“Realtime” dashboards
Echte streaming (subseconde)Directe reactie, continue resultatenComplex, duur, moeilijk te testenFraude, alarmering, live personalisatie
At-least-once + idempotente sinkCorrecte resultaten, betaalbaar, veerkrachtigVraagt gedisciplineerd sleutelontwerpDe meeste streamingpijplijnen
Exactly-oncemachinerieSterke garantie van begin tot eindDuur, beperkt over systemenSmalle paden met hoge inzet
Lambda (batch + snelheid)Nauwkeurige geschiedenis plus vers overzichtTwee codebases te onderhoudenLegacymigraties
Kappa (streaming-first)Eén codebasis, herspeelbaarVraagt bewaard, duurzaam logboekNieuwe streamingplatformen

De centrale spanning is latentie tegenover complexiteit. Elke stap richting realtime kost je operationele last, testmoeilijkheid en geld, en de opbrengsten zijn niet lineair: van dagelijks naar om de paar minuten is goedkoop en vaak genoeg, terwijl van minuten naar subseconde is waar de kosten zich concentreren. Los de spanning op door de beslissing te prijzen, niet de technologie. Vraag welke actie de versheid mogelijk maakt en wat vertraging kost, en koop dan alleen zoveel latentievermindering als die actie rechtvaardigt. Wanneer je streaming wel nodig hebt, leun dan op at-least-oncelevering met idempotente sinks en een streaming-firstlogboek, want die combinatie geeft je correctheid en herspeelbaarheid zonder de zwaarste garanties.

Vragen om met je team te bespreken

  1. Welke beslissing maakt realtime data voor ons werkelijk mogelijk, en wat kost het als die data een minuut te laat aankomt in plaats van direct? Dit is de vraag die elk streamingproject zou moeten poorten, omdat streaming je operationele kosten en complexiteit ruwweg verdubbelt vergeleken met batch. Een groot team kan kwartalen verbranden aan het bouwen van een realtimeplatform dat dashboards bedient die een mens twee keer per dag controleert, wat geld is in brand gestoken. Neem de concrete actie mee die de data drijft, of dat nu een frauduleuze transactie blokkeren is, een operator oproepen of veranderen wat een gebruiker ziet, en zet een getal op de kosten van latentie voor elk. Als het eerlijke antwoord is dat een microbatch van vijf minuten de behoefte zou bedienen, is dat een bevinding om te vieren, niet te verbergen. Het antwoord moet direct veranderen of je echte streaming bouwt, tevreden bent met microbatches of in batch blijft.

  2. Hoe handelen we late en niet-geordende gebeurtenissen af, en wat gebeurt er met data die aankomt nadat een venster sluit? Late en niet-geordende data is het moeilijke deel van streaming, en teams die deze vraag overslaan ontdekken haar in productie wanneer hun getallen weigeren af te stemmen. De concurrerende drukken zijn latentie en correctheid: houd vensters langer open om achterblijvers te vangen en je vertraagt elk resultaat en verbruikt meer geheugen, sluit ze sneller en je laat stilletjes echte data vallen. Neem bewijs mee over hoe laat je data werkelijk aankomt, gemeten als het gat tussen gebeurtenistijd en verwerkingstijd over je bronnen, aangezien een mobiele bron in tunnels heel anders gedraagt dan een gebeurtenis aan de serverkant. Besluit expliciet of late data wordt gelaten vallen, gelogd of een correctie triggert, en zorg dat iedereen stroomafwaarts weet welke. In een overheidscontext waar cijfers verdedigbaar moeten zijn kan het stilletjes laten vallen van late gebeurtenissen een complianceprobleem zijn, dus het beleid moet bewust en gedocumenteerd zijn.

  3. Zijn onze sinks idempotent genoeg om geschiedenis veilig te herspelen, en bewaart ons gebeurtenislogboek genoeg om herspelen mogelijk te maken? Herverwerking is het vermogen dat teams het vaakst wensen te hebben ingebouwd en het vaakst niet bouwden, en het hangt af van twee dingen die samenwerken: idempotente sinks die herspeelde gebeurtenissen absorberen zonder te dupliceren, en een duurzaam logboek dat genoeg geschiedenis bewaart om vanuit te herspelen. Zonder beide betekent een logicabug repareren dat je de getroffen periode niet schoon kunt herberekenen en vastzit aan getallen met de hand patchen onder druk. Neem je huidige bewaartermijn mee en een concrete test: kies een echte bug van het afgelopen kwartaal en vraag of je de gecorrigeerde logica over de getroffen data had kunnen herspelen. De trek ertegen is kosten, aangezien geschiedenis bewaren en idempotente schrijfacties ontwerpen vooraf opslag en discipline vraagt. Maar het alternatief duikt op op het slechtst denkbare moment, tijdens een incident, dus het antwoord bepaalt hoeveel je in herspeelbaarheid investeert voordat je het nodig hebt.

  4. Wanneer een streamingjob crasht, hoe snel moet hij herstellen, hoeveel toestand mag hij vasthouden en hebben we een herstel onder productiebelasting werkelijk getimed? Een batchjob die sterft kan morgen opnieuw draaien, maar een altijd-aanstream die sterft is een uitval die gaande is, en toestandsgebonden jobs die lopende tellingen, joins of fraudemodellen vasthouden kunnen minuten geheugen verliezen of lang nodig hebben om toestand na een herstart opnieuw te laden. Voor een groot team is dit waar een onglamoureus detail stilletjes je echte beschikbaarheid bepaalt: onbegrensde toestand groeit tot een job zonder geheugen raakt, en een trage checkpointherstel verandert een hapering van tien seconden in een van tien minuten. De concurrerende drukken zijn versheid tegenover veiligheid, omdat frequentere checkpoints het herstel verkorten maar overhead toevoegen, en royale toestandsbewaring de nauwkeurigheid verbetert maar geheugenuitputting riskeert. Neem een concreet hersteltijddoel mee, je huidige toestandsgrootte en haar groeicurve, je checkpointinterval en de resultaten van een echte failoveroefening in plaats van een hoopvolle schatting. In omgevingen van onderneming en overheid waar de stroom fraudescoring of een publieke veiligheidsfeed ondersteunt is een niet-getest herstelpad een operationeel risico dat je hebt aanvaard zonder te meten, dus behandel de oefening als eis, niet als aardigheid.

  5. Draaien we één streaming-firstcodebasis of een aparte batchlaag en snelheidslaag, en wat kost het ons werkelijk om de twee afgestemd te houden? Het Lambda-patroon van een batchlaag voor nauwkeurige geschiedenis plus een snelheidslaag voor verse resultaten dwingt je dezelfde bedrijfslogica tweemaal te schrijven, in twee systemen, en hun antwoorden eeuwig af te stemmen, terwijl een streaming-firstvorm (Kappa) een duurzaam, herspeelbaar logboek houdt en alle verwerking als streamverwerking draait. Voor een grote organisatie is de gedupliceerde logica waar afdrijving en omstreden getallen ontstaan, omdat een regel in de ene laag verandert en niet in de andere, en engineers echte tijd besteden aan uitleggen waarom de twee verschillen. De trek om beide te houden is traagheid en het comfort van een bewezen batchlaag, dus weeg dat eerlijk af tegen de onderhoudsbelasting. Neem de lijst berekeningen mee die je nu op beide plekken draait, de incidenten veroorzaakt doordat de twee lagen het oneens waren en een beoordeling of je gebeurtenislogboek genoeg geschiedenis bewaart om batchbehoeften als herspelingen uit te drukken. In overheids- en geauditeerde ondernemingscontexten is het hebben van twee lagen die verschillende cijfers voor dezelfde periode kunnen rapporteren zelf een compliance-verplichting, aangezien je moet kunnen zeggen welk getal gezaghebbend is en waarom.

  6. Wie beheert dit altijd-aansysteem wanneer het om drie uur ‘s nachts breekt, en hebben we de bereikbaarheidslast en de specialistische vaardigheden die het vraagt begroot, of nemen we bezetting naar batchmaat aan? Streaming verschuift kosten van bouwen naar draaien: het systeem moet elke seconde gezond blijven, wat echte bereikbaarheidsdekking betekent, engineers die vloeiend zijn in gebeurtenistijd, watermarks, toestand en leveringssemantiek, en testen dat moeilijker is dan voor een job die draait en stopt. Teams keuren routinematig een streamingplatform goed op grond van zijn vermogens en financieren nooit de mensen die het in leven houden, dus het platform degradeert en het vertrouwen erodeert. De afweging is reikwijdte tegenover houdbaarheid: elke extra realtimepijplijn is nog iets dat iemand kan oproepen, dus de vraag is of de latentie die ze koopt een blijvende operationele verbintenis rechtvaardigt. Neem een eerlijke inventaris mee van wie elke stroom in productie bezit, je huidige bereikbaarheidsrooster en zijn ruimte en waar de gebeurtenistijdexpertise werkelijk zit, of dat nu een aanwerving is, een partner of een beheerde dienst. Voeg voor een publiek orgaan of grote onderneming de doorlooptijden van aanbesteding en werving toe en elke beheerde-dienstoptie, want een realtimeplatform dat afhangt van schaars talent dat je niet kunt werven of behouden is een plan om een uitvalgevoelig systeem onderbezet te draaien.

Sectorperspectief

Startup. Streaming is zelden je eerste zet, en een zwaar platform opzetten kan een klein team laten zinken. Kies het ene signaal dat je kernwaarde raakt, zet gebeurtenissen op één bewaarde op een logboek gebaseerde broker en draai een lichtgewicht processor met op sleutel gebaseerde, idempotente sinks zodat een at-least-onceherpoging nooit dubbel telt. Bewaar een paar dagen geschiedenis zodat je door gerepareerde logica kunt herspelen, en geef de voorkeur aan een beheerde streamingdienst boven je eigen cluster beheren, want je schaarse middel is engineeringaandacht.

Kleinbedrijf. Je hebt waarschijnlijk geen streamingspecialist en geen zin altijd-aaninfrastructuur te draaien, dus behandel realtime als iets wat je koopt binnen tools die je al gebruikt in plaats van een systeem dat je bemant. Formuleer de behoefte als latentievraag met een getal eraan, en in de meeste gevallen voldoet een microbatch die om de paar minuten ververst tegen een fractie van de kosten en het risico. Kies leveranciers wier realtimefuncties transparant zijn over vertraging en makkelijk terug te vallen, en reserveer maatwerkstreaming voor het zeldzame geval waar verse data direct omzet of veiligheid drijft.

Grote onderneming. Het probleem is consistentie en kosten over veel teams: een gedeeld op logboek gebaseerd platform, een standaardbeleid voor gebeurtenistijd en late data en idempotente sinks zodat groepen ophouden kwetsbare pijplijnen opnieuw uit te vinden. Begroot de altijd-aanoperaties en bereikbaarheidslast expliciet, standaardiseer op een streaming-firstlogboek zodat je een gedupliceerde batchcodebasis vermijdt en beheer streams als bestuurde dataproducten met eigenaren, schemaversiebeheer en bewaakte vertraging in plaats van een verstrooiing van maatwerkjobs. Volg latentie, hersteltijd en kosten per stroom als portfoliostatistieken.

Overheid. Controleerbaarheid en publieke verantwoording geven elke keuze vorm. Bewaar elke verwerkte gebeurtenis in een duurzaam logboek zodat cijfers gerapporteerd aan toezichtsorganen, reizigersaantallen, anomalieën in uitkeringen, fraudebeslissingen, exact kunnen worden gereconstrueerd, en maak het beleid voor late data expliciet en gedocumenteerd in plaats van gebeurtenissen stilletjes te laten vallen. Aanbesteding moet overdraagbaarheid van data eisen en bekendmaking van de lever- en bewaargaranties van een beheerde dienst, en elke herziening na een regelwijziging moet een verdedigbare herspeling door gecorrigeerde logica zijn, geen handmatige patch die niemand kan traceren.

Voorbeelden

Startup. Een consumentenapp wil gebruikers een live activiteitenfeed tonen en verdachte logins markeren terwijl ze gebeuren. Het team weerstaat een zwaar streamingplatform opzetten. Ze zetten gebeurtenissen op één bewaarde op een logboek gebaseerde broker, draaien een lichtgewicht streamprocessor voor de loginrisicologica en voeden een realtime-OLAP-opslag die de activiteitenfeed aandrijft. Elke sink is op sleutel gebaseerd en idempotent, dus een at-least-onceherpoging telt nooit dubbel. Wanneer ze later een bug in de risicoregel vinden, spelen ze het logboek ‘s nachts simpelweg door de gerepareerde logica opnieuw af, omdat ze een week geschiedenis bewaarden en nooit een tweede batchcodebasis nodig hadden.

Grote onderneming. Een retailbank scoort elke kaarttransactie op fraude binnen het autorisatievenster, de live transactiestroom joinend tegen een toestandsgebonden model van recent rekeninggedrag. Checkpointing laat de scoringsdienst in seconden herstellen van een knooppuntfalen zonder zijn geheugen van de laatste paar minuten te verliezen. Daarnaast streamt change data capture updates uit de kernbankingdatabase naar een zoekindex en een personalisatiedienst, beide vers houdend zonder pollen. Operationele dashboards lezen uit een realtime-OLAP-opslag zodat risico- en operatieteams het bedrijf zien bewegen, en de hele pijplijn zendt de lag- en doorvoertelemetrie uit beschreven in hoofdstuk 9.2.

Overheid. Een vervoersautoriteit van een grote stad neemt voertuigposities en vervoersbewijstikken in om aankomsten te voorspellen en drukte realtime te bewaken, voor zowel publieke apps als een operatiecentrum. Omdat reizigers in tunnels tikken in vertraagde bursts uploaden, berekent het team reizigersaantallen op gebeurtenistijd met watermarks afgestemd op de waargenomen vertraging, en logt het elke gebeurtenis die aankomt nadat haar venster sloot in plaats van haar stilletjes te laten vallen. Elke verwerkte gebeurtenis wordt bewaard in een controleerbaar logboek zodat reizigersaantallen gerapporteerd aan toezichtsorganen exact kunnen worden gereconstrueerd. Wanneer een tariefregel verandert, spelen ze de getroffen periode door de gecorrigeerde logica opnieuw af en produceren een verdedigbare herziening.

Zakelijke onderbouwing: motivatie, ROI en TCO

Het rendement van realtimedata komt uit handelen terwijl handelen nog telt. Fraude gevangen tijdens autorisatie voorkomt een verlies dat een nachtelijke batch alleen zou rapporteren. Personalisatie die binnen een sessie reageert verhoogt conversie op een manier die de aanbeveling van morgen niet kan. Operationele bewaking die het heden weerspiegelt laat je ingrijpen voordat een klein probleem een uitval of publiek incident wordt. In elk geval is de waarde het verschil tussen nu handelen en later handelen, en dat verschil is wat je moet kwantificeren wanneer je de zaak maakt.

De total cost of ownership is hoger dan batch, en eerlijk zijn daarover beschermt je geloofwaardigheid. Je betaalt voor altijd-aaninfrastructuur, voor engineers die gebeurtenistijd, watermarks, toestand en leveringssemantiek begrijpen en voor de moeilijkere test- en bereikbaarheidslast van een systeem dat continu gezond moet blijven in plaats van te draaien en te stoppen. Een streaming-firstarchitectuur op een bewaard logboek verlaagt de doorlopende kosten door je een dubbele batchcodebasis te besparen, en at-least-once met idempotente sinks kiezen vermijdt de kosten van exactly-oncemachinerie van begin tot eind. De duurste fout is realtime bouwen waar microbatch of batch volstond, dus het sterkste kostenargument is vaak een beslissing om niet te streamen. Formuleer het verhaal voor het bestuur rond specifieke latentiegevoelige beslissingen en hun meetbare opbrengst, en wees even duidelijk over waar in batch blijven geld bespaart zonder waardeverlies.

Antipatronen en valkuilen

  • Streaming bouwen voor prestige terwijl een microbatch om de paar minuten aan de behoefte voldeed.
  • Op verwerkingstijd rekenen, zodat je getallen met je infrastructuur wiebelen in plaats van met de wereld.
  • Late en niet-geordende data negeren tot afstemming in productie faalt.
  • Letterlijk exactly-once overal najagen in plaats van at-least-once met idempotente sinks.
  • Onbegrensde toestand zonder verloop, die stilletjes groeit tot een job zonder geheugen raakt.
  • Geen checkpointing, zodat een herstart toestand verliest of een volledige herspeling afdwingt.
  • Operationele databases op een timer pollen in plaats van change data capture te gebruiken.
  • Een Lambda-batchlaag en snelheidslaag onderhouden met gedupliceerde, afdrijvende logica.
  • Een bewaartermijn te kort om geschiedenis te herspelen wanneer je een bug vindt.
  • Streams zonder lag-, doorvoer- of versheidsstatistieken, die stilletjes falen.

Volwassenheidsmodel

  • Niveau 1, Initiëren: Alles is batch, of een paar met de hand gebouwde streamingjobs draaien reactief zonder bewaking. Getallen worden berekend op verwerkingstijd, late data wordt genegeerd en een herstart verliest toestand. Niemand kan geschiedenis herspelen om een bug te repareren, en problemen worden ontdekt wanneer cijfers stroomafwaarts weigeren af te stemmen.
  • Niveau 2, Ontwikkelen: Sommige teams draaien kernstreamingpijplijnen op een op logboek gebaseerde broker met checkpointing, en ze onderscheiden gebeurtenistijd van verwerkingstijd en gebruiken basisvensters. De praktijk is inconsistent van team tot team: levering is at-least-once maar niet alle sinks zijn idempotent, afhandeling van late data is geïmproviseerd en lag wordt informeel bekeken in plaats van gealarmeerd.
  • Niveau 3, Standaardiseren: Gebeurtenistijd, watermarks en een expliciet beleid voor late data zijn gedocumenteerd en over de organisatie toegepast. Sinks zijn idempotent voor effectief-eenmalige resultaten, toestand heeft verloop en change data capture voedt downstreamsystemen bij conventie. Een bewaard logboek ondersteunt herspelen, en lag, doorvoer en versheid worden bewaakt met alarmen als organisatiebrede standaard in plaats van gewoonte per team.
  • Niveau 4, Beheersen: Het streaminglandschap wordt gemeten en beheerst aan de hand van uitgangswaarden. Elke pijplijn draagt service-level objectives voor latentie van begin tot eind, afnemersvertraging, hersteltijd, afwijking van gebeurtenistijd, percentage late gebeurtenissen, toestandsgrootte en kosten per miljoen gebeurtenissen, allemaal gevolgd tegen afgesproken doelen en alarmerend bij regressie. Herstel wordt geoefend en getimed in plaats van aangenomen, tegendrukruimte en toestandsgroei worden bewaakt als capaciteitssignalen en een nieuwe stroom moet deze statistieken halen voordat hij naar productie gaat.
  • Niveau 5, Orkestreren: Een streaming-firstarchitectuur bedient zowel verse als historische behoeften uit één herspeelbaar logboek, en streaming-SQL, gematerialiseerde views en realtime-OLAP maken verse data breed toegankelijk. Herverwerking is routine en getest, het platform schaalt automatisch en herbalanceert tegen gemeten belasting en kosten, en streams worden afgeschaft, opnieuw afgebakend of vervangen op bewijs. Streaming is geïntegreerd met bedrijfs- en risicoplanning, en elke stroom is van begin tot eind observeerbaar en controleerbaar naarmate het belasting- en kostenbeeld verschuift.

Ideeën voor discussie

  1. Waar in je stack verdient “realtime” werkelijk zijn kosten, en waar is het een onbezonnen wens?
  2. Hoe groot is het gat tussen gebeurtenistijd en verwerkingstijd over je bronnen, en meet je het?
  3. Kun je een Lambda-opzet met batch en snelheid samenklappen tot één streaming-firstcodebasis, en wat zou dat blokkeren?
  4. Welke van je sinks zijn werkelijk idempotent, en kon je de data van vorig kwartaal vandaag veilig door gecorrigeerde logica herspelen?
  5. Wat is je beleid voor data die aankomt nadat een venster sluit, en weet iedereen stroomafwaarts dat?
  6. Hoe zou change data capture veranderen hoe je zoeken, caches en analytics synchroon houdt?

Belangrijkste inzichten

  • Grijp alleen naar streaming wanneer een latentiegevoelige beslissing het betaalt. Batch en microbatch zijn goedkopere standaarden.
  • Reken op gebeurtenistijd en behandel late en niet-geordende data als kernprobleem, afgehandeld met vensters en watermarks.
  • Geef de voorkeur aan at-least-oncelevering met idempotente sinks voor effectief-eenmalige resultaten boven letterlijk exactly-once overal.
  • Checkpoint toestandsgebonden verwerking, begrens je toestand en bewaak afnemersvertraging als kopstatistiek.
  • Gebruik change data capture om vanuit operationele databases te streamen in plaats van te pollen.
  • Geef de voorkeur aan een streaming-firstarchitectuur op een bewaard, herspeelbaar logboek boven twee codebases onderhouden.
  • Stel streams beschikbaar via streaming-SQL, gematerialiseerde views en realtime-OLAP, en houd elke stroom observeerbaar en controleerbaar.

Referenties en verder lezen

  • Tyler Akidau, Slava Chernyak, and Reuven Lax, “Streaming Systems.”
  • Martin Kleppmann, “Designing Data-Intensive Applications.”
  • Nathan Marz and James Warren, “Big Data” (Lambda architecture).
  • Jay Kreps, “Questioning the Lambda Architecture” (O’Reilly Radar).
  • Fabian Hueske and Vasiliki Kalavri, “Stream Processing with Apache Flink.”
  • Ben Stopford, “Designing Event-Driven Systems.”
  • Tyler Akidau and colleagues, “The Dataflow Model” (VLDB paper on windowing and watermarks).