SpecsModelExplore
Overview
Actors
External Systems
IoT Platform
Edge
Ingestion
Stream Processing
Stream Processor
Kafka ConsumerWindow AggregatorAnomaly Detector
Batch Processor
Device Management
User-Facing Apps
Platform
Kafka Consumer
How it reads
Deserialize
Forward
Notes
SpecsModelExplore
Kafka Consumer
OverviewGuideLinks and Communications

Kafka Consumer

The Kafka Consumer is the entry point of the streaming pipeline. It reads the validated telemetry topics published by ingestion and hands records to the Window Aggregator with exactly-once semantics.

Offsets and checkpoints

Consumer offsets are committed as part of Flink checkpoints, not independently. On recovery the pipeline rewinds to the last checkpoint so no window is double-counted and none is lost.

How it reads

Subscribe

The consumer subscribes to the telemetry topics, spread across partitions keyed by device id so each device’s events stay ordered.

Deserialize

Records are decoded and lightly validated; malformed payloads are routed to a dead-letter topic instead of stalling the pipeline.

Forward

Valid events are forwarded downstream, preserving event-time timestamps for correct windowing.

Notes

  • Partitioning by device id keeps per-device ordering and enables parallel consumption across the configured task slots.
  • Back-pressure from downstream operators naturally throttles consumption rather than dropping events.