Clarify requirements & constraints
- Millions of devices, high throughput (100k–M msgs/sec), mixed protocols (MQTT/HTTP/CoAP), per-device ordering, dedup, low-latency analytics + batch historical processing, multi-tenant, secure provisioning.
High-level architecture
- Edge gateways for constrained devices; global API Gateway / Load Balancer -> Ingestion Tier -> Streaming Platform -> Stream processors + Batch pipeline -> Storage (hot/cold) -> Downstream services/analytics.
Ingestion & scaling
- Use protocol-optimized front doors (managed IoT service like AWS IoT Core / Azure IoT Hub or API GW + MQTT brokers). Autoscale frontends behind global LB and regionally deployed ingress clusters.
- Push to a partitioned durable message bus: Apache Kafka (MSK) or cloud equivalents (Kinesis, Event Hubs). Partition by device_id or device_shard to preserve ordering scope.
Message ordering & deduplication
- Ordering: guarantee ordering per partition (use device_id -> partition key or shard-per-device-group). For strict single-device ordering, use single-partition stream per device-group.
- Deduplication: include monotonic sequence numbers + device-generated UUID or HMACed message id. Implement idempotency at streaming ingress using a dedupe store with TTL (Redis streams / DynamoDB with conditional writes). For high throughput, use compacted Kafka topic of seen-ids per shard or Bloom filters + verification in downstream processors to reduce state.
Downstream processing
- Real-time: Stateful stream processors (Apache Flink, Kafka Streams, or cloud Kinesis Data Analytics) for anomaly detection, enrichment, immediate alerts. Maintain keyed state per device (windowing, late arrivals).
- Batch: Periodic ETL with Spark on data lake for historical analytics, ML training. Stream -> nearline store/object lake (Parquet on S3/ADLS) for batch access.
- Use change-data capture and Lambda architecture patterns where necessary; consider stream-first to minimize latency.
Storage choices
- Hot time-series: scalable TSDB (InfluxDB/TimescaleDB or Cassandra/Scylla for wide rows) or cloud managed time-series (Timestream) for recent queries.
- Cold/analytics: Object storage (S3/ADLS) in partitioned Parquet, cataloged via Glue/Databricks.
- Metadata and device registry: strongly-consistent KV (DynamoDB/CosmosDB) with TTL for ephemeral keys.
- Long-term dedupe/idempotency: compacted Kafka topics or DynamoDB with conditional writes.
Security & device-to-cloud at scale
- Device identity: per-device X.509 certificates or hardware-backed keys (TPM/IoT Secure Element). Use device provisioning service (DPS) to issue certs, manage attestation and ownership.
- Transport: mTLS for MQTT/HTTPS; enforce TLS1.2+. Mutual auth to prevent spoofing.
- Auth & authorization: short-lived JWTs issued by provisioning/orchestration service, scope-limited roles, RBAC for tenants.
- Key rotation & revocation: centralized PKI, CRL/OCSP or certificate revocation via registry. Support device re-provisioning and onboarding flows.
- Network: VPC endpoints, private links to avoid public internet for ingestion to backend.
- Rate limiting, quota per device/tenant, anomaly detection for compromised devices.
- Observability & auditing: centralized logging, metrics, SIEM integration, per-device telemetry policy.
Scalability & trade-offs
- Partitioning controls ordering vs parallelism; too-fine partitions increase coordination overhead.
- Stateful stream processors require careful state sharding and checkpointing.
- Dedup store TTL balances memory vs safety window.
- Managed services reduce ops overhead; self-managed gives cost/control.
Outcome: a regionally distributed, partitioned stream-first architecture with strong device identity, dedupe and ordering via partition keys + idempotency store, real-time processing in Flink/KStreams and long-term analytics in a data lake—scalable, secure, and operable.