Designing a fault-tolerant Industrial IoT data pipeline โ from raw sensor MQTT events to hourly analytics โ with a self-healing PostgreSQL HA cluster at its core.
RoleCloud & Data Architect
PlatformAWS EC2 ยท ap-south-1
DomainIndustrial IoT / Manufacturing
TypePOC โ Production Pipeline
The Problem
Sensor data arriving faster than a single database could handle โ with zero tolerance for data loss
Industrial sensors emit telemetry every few seconds. The existing setup wrote directly to a single PostgreSQL instance with no replication, no buffering, and no failover. Any database restart meant lost telemetry during the window.
Core Requirement: Zero data loss during planned or unplanned database downtime. Automated failover under 30 seconds. Hourly IoT aggregations delivered to the analytics layer on schedule.
Architecture
Three-tier pipeline: ingest โ buffer โ store โ analyse
1
MQTT Broker (Mosquitto)Sensors publish telemetry to Mosquitto. Python subscribers consume events and forward to Kafka โ keeping sensor firmware simple with no Kafka client on device.
2
Apache Kafka (Docker on EC2)Kafka acts as a durable buffer. Events are retained in partitions even if the downstream DB is unavailable, and consumed once it recovers. Topic: iot-telemetry, 3 partitions, 7-day retention.
3
PostgreSQL 17 HA Cluster (Patroni + etcd + HAProxy)Three-node cluster: 1 primary + 2 replicas. Patroni manages leader election via etcd. HAProxy routes writes to primary and reads to replicas. Automatic failover under 30 seconds.
4
pgBackRest S3 Backups with Glacier TieringContinuous WAL archiving to S3. Daily full backups with 30-day retention tiered to Glacier after 7 days. RPO under 5 minutes.
5
Apache Airflow + PySpark AnalyticsHourly DAG triggers a PySpark job computing min/max/avg aggregations per sensor per hour, written to the analytics schema for dashboards.
Key Challenges
What nearly derailed the project
โก
Split-brain risk during network partitionWithout strict etcd quorum configuration, a network blip could cause both nodes to claim primary. Fixed by enforcing majority quorum and adding DCS timeout before any failover begins.
๐พ
WAL slot exhaustion filling EC2 diskDuring a Kafka consumer outage, unconsumed WAL slots grew 18GB in 4 hours. Fixed with max_slot_wal_keep_size and CloudWatch alarm on WAL directory size.
โฑ๏ธ
Airflow DAG timing driftPySpark jobs occasionally exceeded 60 minutes causing queue pile-up. Resolved with task timeout enforcement and maximum active runs limit of 1.
๐
HAProxy false-positive health checksHAProxy marked the primary unhealthy during a long VACUUM, triggering unnecessary failover. Fixed by using Patroni's REST API health endpoint instead of TCP connection check.
Outcomes
A pipeline that runs without hand-holding
<30s
Automated DB failover time
0
Data loss events post-launch
7d
Kafka retention buffer
<5min
Recovery point objective
Result: Millions of sensor events processed with no data loss incidents. Two automated failovers completed (planned maintenance) with zero manual intervention.