Direct naar inhoud
Terug naar projecten
Kafka
Streaming
Python
Postgres
Docker

Realtime OV-streamingpipeline op Kafka

Eigen project

3 weken

Eigen project, gebouwd omdat streaming precies het terrein is waar je een data engineer op beoordeelt: windowing, idempotentie, foutafhandeling — en of het blijft draaien als niemand kijkt.

Wat het is

Een streamingpipeline die de realtime voertuigposities en aankomstvoorspellingen van het Nederlandse openbaar vervoer (NDOV/OVapi) verwerkt via Kafka, per lijn en per station voortschrijdende vertragingsstatistieken berekent in 5-minuten-windows, en die live toont op een kaart: waar rijdt alles nu, en hoe erg is de vertraging.

De lat: de repo moet lezen als productiewerk, niet als een tutorial. Dus tests, CI, een dead-letter queue, observability — en een README die de vragen beantwoordt die er echt toe doen bij streaming: doorvoer, lag-gedrag en schema-evolutie.

De architectuur

NDOV / OVapi feed
      │  producer (Python)

 Kafka-topics (Redpanda)
      ├──► raw archief (MinIO) — elke ruwe message, herafspeelbaar

 stream processor — 5-min-windows per lijn / station
      ├──► DLQ-topic (onverwerkbare messages)

 Postgres (aggregaten) ──► live dashboard met vertragingskaart
      ·
 Prometheus + Grafana: lag, doorvoer, DLQ-diepte

Alles draait lokaal met één commando (make up), zonder cloudkosten. Redpanda als Kafka-compatibele broker, omdat het licht is en zich exact als Kafka gedraagt.

De ontwerpkeuzes die het verschil maken

  • Idempotente verwerking. De consumer commit zijn offset pas ná de upsert in Postgres, en de upsert-sleutel is (window, dimensie). Herlevering kan een afgesloten window dus nooit dubbel tellen — de standaardvraag over at-least-once-verwerking, beantwoord in code én test.
  • Dead-letter queue. Onverwerkbare messages zijn bij publieke feeds een zekerheid, geen uitzondering. Ze landen met foutmetadata in een eigen topic, zijn zichtbaar in Grafana, en na een parserfix worden ze met één commando opnieuw verwerkt.
  • Ruw archief vóór verwerking. Elke message wordt onbewerkt gearchiveerd voordat er iets mee gebeurt. Gaat er iets mis in de verwerking, dan is er altijd een herafspeelpad vanaf de bron.
  • Schema-evolutie. Events dragen een schemaversie; de parser ondersteunt de huidige en de vorige. Hoe een volgende versie uitrolt zonder de pipeline te breken staat gedocumenteerd.
  • Late events. Een vaste tolerantie vóór het afsluiten van een window; wat later binnenkomt wordt geteld als metriek en bewust losgelaten. Een afweging die je documenteert, niet wegmoffelt.

De meetlat — en het resultaat

De lat lag vooraf vast en is gehaald: een schone clone start met make up de volledige stack, en de pipeline heeft een aaneengesloten run van 24 uur doorstaan op het landelijke feedvolume — zonder ingrijpen, met een consumer-lag die begrensd bleef in plaats van gestaag op te lopen. De gemeten doorvoer staat in de README, naast de Grafana-grafieken van de run en een schermopname van het live dashboard. Een test bewijst dat dubbele levering de cijfers niet vervuilt, en de DLQ-route is gedemonstreerd van fout bericht tot herverwerking.

Waarom dit ernaast staat

Het meeste dagelijkse datawerk is geen streaming — dashboards, warehouses en rapportages hebben dat zelden nodig. Maar de discipline is dezelfde: fouten expliciet afhandelen, verwerking idempotent maken, alles meetbaar. Dit project laat zien dat die aanpak overeind blijft als de data niet één keer per nacht komt, maar duizenden keren per minuut.

In de repo

  • Producer, ruw archief en Kafka-backbone
  • Windowed processor met DLQ en idempotente sink
  • Live dashboard met vertragingskaart
  • Observability: Prometheus en Grafana met dashboards voor lag, doorvoer en DLQ-diepte
  • Een duurtest van 24 uur op het landelijke feedvolume; gemeten doorvoer en schermopname in de README

Technologie

Kafka (Redpanda), Python, Postgres, MinIO, Streamlit, Prometheus, Grafana, Docker Compose