Q: Your company deploys smart grid sensors across 10 million industrial devices. Every device connects over MQTT and transmits sensor telemetry every 30 seconds, generating 333,000 incoming messages per second. Devices operate over unstable cellular networks with frequent disconnects. How do you design an ultra-scalable, resilient IoT ingestion pipeline that guarantees message deduplication and sub-second anomaly querying?
Architectural design for ingesting, validating, and querying real-time telemetry from 10,000,000 concurrent IoT devices (100k msgs/sec sustained, 1M msgs/sec peak) using EMQX MQTT brokers, Kafka, and TimescaleDB.
Want to master this scenario in a live sandbox? The Linux Foundation's FinOps Certified Practitioner (FOCP) Program covers this exact problem with hands-on terminal drills.
🛠️ Production Runbook & Step-by-Step Resolution
Deploy Horizontally Scaled EMQX Distributed MQTT Broker Fleet
Terminate millions of persistent lightweight MQTT connections at the edge:
- EMQX Cluster: Deployed a 16-node EMQX enterprise cluster using Erlang OTP lightweight processes, capable of maintaining 10+ million concurrent persistent MQTT TCP/TLS connections.
- Keep-Alive & Session Expiry: Configured MQTT Keep-Alive (60s) and Clean Start session handling to gracefully resume connections over flakey cellular links without dropping queued messages.
Bridge Telemetry from MQTT into Partitioned Apache Kafka Topics
Decouple real-time device message ingestion from backend database processing:
- EMQX Rule Engine: Configured native EMQX high-throughput rule engine to stream incoming MQTT payloads directly into Kafka topics:
iot.sensor.telemetry. - Partition Key: Partitioned Kafka topics by
device_idto guarantee strict in-order message delivery per individual sensor.
Process Streams & Deduplicate with Apache Flink Stream Processing
Filter duplicate messages and calculate real-time moving averages:
- Deduplication: Flink maintains a 10-minute stateful window evaluating
(device_id, message_seq_num)to filter duplicate retransmissions caused by cellular network blips. - Anomaly Detection: Flink computes 5-minute rolling averages of temperature and voltage; if values spike >3 standard deviations, an immediate alert is published to an emergency Kafka topic.
Store Time-Series in TimescaleDB with Automated Hypertables & Compression
Optimize multi-billion row relational time-series storage and analytical queries:
- Hypertable Partitioning: Created TimescaleDB hypertables partitioned into 1-day chunks by timestamp.
- Columnar Compression: Automatically compressed chunks older than 2 days:
ALTER TABLE sensor_data SET (timescaledb.compress, timescaledb.compress_segmentby = 'device_id'), achieving a 92% storage footprint reduction. - Query Performance: Analytical queries across 10 billion data points execute in < 180ms.
- Terminate 10M concurrent persistent IoT connections using a distributed EMQX MQTT cluster.
- Stream MQTT payloads into Apache Kafka partitioned by device_id for ordered buffering.
- Deduplicate and detect anomalies in real time using Apache Flink stateful stream processing.
- Store billions of sensor readings in TimescaleDB hypertables with 92% columnar compression.