SpecsModelExploreRequirements
SpecsModelExploreRequirements
Kafka Consumer
Overview
Guide
Twin streams
Links and Communications
IDproduct.processing.stream_processor.kafka_consumer
DescriptionConsumes telemetry topics.
Keykafka_consumer
Typecomponent
Statusactive
TeamStream Processing
Owner
stream-processing@iotgateway.example
TechnologiesJava, Kafka
Tags
Table of contents
Guide
Kafka Consumer
How it reads
Deserialize
Forward
Notes
Links and Communications
Requirements

Attributes

partitions128
throughput500k events/s

Guide

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.

Links and Communications

1 incoming link

Incoming Links(1)

Source objectLabelLink typeRelationRaw idRequiredDescription
Window Aggregator

example-iot.product.processing.stream_processor.window_aggregator

component
aggregates stream from—linksproduct.processing.stream_processor.kafka_consumer——

Requirements