← All writing
articleJan 09, 202519 min read

IoT Data Ingestion: The Platform Between a Device and Kafka

Connection counts, QoS, session state, provisioning, and the MQTT to Kafka bridge, written for fleets large enough that the hard problem is connections rather than throughput.

MQTTKafkaIoTArchitecture
IoT Data Ingestion: The Platform Between a Device and Kafka cover illustration

Devices do not speak Kafka. They speak MQTT, CoAP, or a vendor protocol over a constrained network, and something has to sit between that world and the event log where the rest of the system lives.

That something is usually described as an ingestion gateway. It is also, quietly, the component that decides whether the rest of the platform is reliable, because every guarantee you believe you have is only as good as this boundary.

This article is about that boundary at a scale where the interesting problems are not throughput.

The fleet and what it does to a broker

Assume a consumer IoT platform with five million devices, each publishing telemetry every thirty seconds.

5,000,000 devices / 30 sec = 166,667 messages/sec average

With reconnects, upgrades, and bursts, design for roughly ten times that: about 1.7 million messages per second at peak.

The daily volume follows.

166,667/sec x 86,400 sec     = 14.4 billion messages per day
at ~200 bytes per message    = 2.88 TB per day
                               = roughly 1 PB per year

One petabyte a year is a storage and cost conversation. It is not the hard part.

The hard part is this: those five million devices are online at the same time, and each one is a TCP connection with a TLS session, a keep-alive timer, a subscription set, and protocol state.

5,000,000 concurrent connections
at ~50,000 connections per broker node
                              = 100 broker nodes

A hundred nodes is a real, operable, autoscaling fleet. It is also the number that should drive your capacity plan, and most ingestion designs start by planning the message rate instead.

Connections, not messages, are the constraint

A broker node is comfortable with 166,667 divided by a hundred messages per second. That is roughly 1,667 messages per second per node, which is a rounding error.

The same node has to hold 50,000 connections. Per connection the broker tracks at minimum:

  • the socket and the TLS state, which is the largest single item;
  • the client’s identifier and keep-alive deadline;
  • the subscription set;
  • the in-flight QoS 1 and QoS 2 message identifiers and their delivery state;
  • ping timers and any last-will message.

Multiply by 50,000 and the memory profile becomes the design problem. A broker configured happily for 10,000 connections will start evicting sessions long before it runs out of message throughput, and the symptom is a device that appears to be connected but receives nothing.

The three numbers to measure per node, in this order:

  1. Connection count and connection churn per second. Churn matters more than the total. A thousand stable connections are easier than fifty connections reconnecting ten times a second.
  2. Resident memory per connection. Compare it against your node size to find the real ceiling, rather than trusting the configured limit.
  3. Session store size if sessions are persisted anywhere, covered next.

MQTT is not Kafka and the difference matters

MQTT is a client-to-server publish and subscribe protocol designed for constrained devices and unreliable networks. Kafka is a partitioned, replicated, replayable log. They share a shape and almost nothing else.

MQTT Kafka
Unit of ordering per client connection per partition
Replay no, only undelivered while the session lives yes, from any offset
Retention as long as the session persists time or size based, independent of consumers
Consumer identity client id consumer group
Typical fan-out broker to many subscribers on a topic consumer group reads once
Delivery QoS 0, 1, or 2 per message at-least-once, or transactions
Backpressure signal via a receive maximum lag, with the log absorbing bursts

The row that causes the most trouble is ordering. It is the first.

Ordering stops at the connection

MQTT guarantees ordering for messages published by one client over one connection, in the order that client published them.

Two devices publishing to the same topic get no ordering relationship relative to each other, and MQTT says nothing about how a broker interleaves them across subscribers. A naive bridge to Kafka inherits exactly that ambiguity.

The fix is in the Kafka message key, not the protocol. Key on the thing whose order you actually care about.

wrong:  key = topic name
        -> all devices on "temp/kitchen" interleave

better: key = deviceId + ":" + metric
        -> one device's temperature stream is ordered
        -> one device's humidity stream is ordered independently

Once keyed that way, per-device and per-metric order is a property of the partition assignment. You gave it up at the protocol boundary and got it back for free in the log, provided you chose the key deliberately rather than letting a library pick a round-robin.

QoS: three promises you have to implement

MQTT’s quality of service levels are a per-message delivery promise, and each one costs something on a device that may be on a battery.

QoS 0, at most once. Send and forget. No acknowledgement, no redelivery. A message lost in transit is gone.

QoS 1, at least once. The broker acknowledges receipt and may redeliver until it is acknowledged. Every consumer must be safe under repetition.

QoS 2, exactly once. A four-step handshake between client and broker. Expensive in round trips and in broker state, and it only guarantees the hop between that client and that broker. It says nothing about what happens after the message leaves for Kafka.

That last point deserves emphasis, because it is the most common misunderstanding in this area. QoS 2 does not give you exactly-once delivery into your data platform. It gives you exactly-once between one device and one broker. The moment a message is bridged to a topic, duplicated across partitions, or replayed from a consumer group, that guarantee is gone.

For most telemetry, QoS 0 is the correct choice. A missing temperature reading is replaced thirty seconds later. Paying for acknowledgements and session state on every message to protect a number that is superseded is a poor trade.

Reserve QoS 1 for events where a loss is not acceptable: a door opening, a payment terminal reporting, a safety interlock changing state.

Persistent sessions are state you have to store

An MQTT session is persistent when the client connects with a clean start flag disabled. The broker then remembers the subscriptions and queues QoS 1 and 2 messages for delivery when the client returns.

This is a genuinely useful feature and it is where the arithmetic turns frightening.

5,000,000 devices
one message every 30 sec
                     = 2,880 messages per device per day
5,000,000 x 2,880   = 14.4 billion queued messages
at 200 bytes        = 2.88 TB of broker-side queue

To buffer a single day for a fully disconnected fleet, your broker tier needs three terabytes of durable queue, plus the session index, plus the headroom for the messages arriving during recovery.

And the 2.88 terabytes is not a rare condition. It is what a fourteen-hour network outage looks like.

So the design conclusion is forced, and it should be stated as policy rather than discovered:

high-frequency telemetry  ->  QoS 0, clean session, no retention
commands and critical events -> QoS 1, persistent session, short session expiry

Session expiry is the lever. MQTT 5 lets a client or server set a session expiry interval, and a session expiry of zero means the session ends at disconnect, which discards anything queued. Setting a short expiry, minutes rather than hours, bounds the storage you can be held to without giving up the feature for brief outages.

A broker cluster designed for QoS 1 and persistent sessions across five million devices is a different, much larger system than one designed for QoS 0 with clean sessions. Both are buildable. Deciding which one you are building before you provision anything is the whole point.

Provisioning is the real auth design

Five million devices need five million credentials, and the credential story usually gets less attention than the message path even though it is where breaches happen.

The options:

A shared fleet username and password. One credential for all devices. Revocation is all or nothing, an attacker with it can impersonate any device, and it appears in firmware on five million units. It is not acceptable at this scale, and at small scale it is merely bad.

Per-device username and password in a hashed store. Workable, and it means a credential store with five million rows, a provisioning pipeline that writes to it, and a device that can be individually revoked. It also means a hundred thousand passwords sitting in device flash.

A pre-shared key per device. Better than a shared fleet key because it is individually revocable and there is nothing to crack. It is symmetric, so the device can impersonate the server, which matters for commands.

A client certificate per device. Mutual TLS. The device proves possession of a private key held in a secure element, the broker proves its own, and neither side presents a bearer token. Individually revocable, and the credential is not extractable from the device.

For a fleet this size, mutual TLS is the right default, and the work is entirely in the certificate lifecycle.

5,000,000 certs, 90-day lifetime
                            = 55,556 issued per day

Fifty-five thousand issuances a day is a real PKI workload, and it brings the revocation question with it. Do you publish a CRL, or use OCSP, or avoid revocation entirely?

The alternative that removes the question is short-lived certificates, where a device with a hardware root of trust fetches a fresh 24-hour certificate on boot using a signed attestation. Nothing to revoke, because nothing lives long enough to be worth revoking, and a stolen key is useful for less than a day.

The cost is issuance volume, which goes up by the ratio of the lifetimes. Ninety-day certs at 55,556 a day becomes twenty-four-hour certs at five million a day. That is a decision about your PKI’s throughput, and it is a real one rather than a free security win.

Whichever you choose, the provisioning flow needs a way to get the first credential onto a device that has no credentials, and that bootstrap step is where the security actually lives. A device that will accept a certificate from anything that can reach it is not authenticated, whatever happens afterwards.

Topic design and the wildcard trap

Topics carry structure, and the structure is an API you cannot change cheaply.

vending/{tenantId}/device/{deviceId}/telemetry/{metric}

This is reasonable and it creates a real problem. If a device may publish only its own subtree, the broker has to decide whether vending/t-1/device/d-99/telemetry/temp is permitted for a client authenticated as d-42. That is a pattern match against a topic containing another device’s identity, evaluated on every publish, with millions of rules in the store.

The common resolution, and it is a genuine trade-off rather than a compromise:

  • enforce the prefix at the client, where the device only ever constructs topics for itself;
  • enforce the tenant boundary at the broker, which is a small number of rules;
  • rely on the connection identity to identify the device, and treat the topic’s device id as untrusted metadata to be validated by the bridge.

That last point is the one people get wrong. If the bridge trusts the device id in the topic, then any device that can publish can publish as another device, and your data is untrustworthy regardless of how good the TLS was.

Wildcard subscriptions compound this. A client subscribing to vending/+/device/+/telemetry/# matches a very large number of topics, and each match costs broker memory and per-publish work. Per-device wildcards are fine. Fleet-wide wildcards are a capacity problem, and they are usually a sign that the architecture wanted a stream processor and got a broker subscription instead.

The bridge from MQTT to Kafka

The bridge subscribes to the broker and produces to Kafka. It is the least interesting component on paper and the most reliable source of subtle bugs.

The decisions:

Grouping. Publishing one Kafka record per MQTT message at 166,667 per second is wasteful. Batch by topic and by key over a short window, a few hundred messages or a few hundred milliseconds, and record the batch’s receive-time range in the header. Batching is the difference between a few hundred MB per second of broker egress and tens of MB.

What to use as the key. Already covered: device plus metric.

Backpressure. If Kafka is unavailable, the bridge stops consuming. MQTT then applies its own flow control, and the broker’s queues fill. This is where the 2.88 terabytes from earlier becomes real. A bridge must either spill to local disk with a declared capacity limit, or drop, and it must be explicit which. A silently truncated stream is the worst outcome because it looks healthy.

Deduplication. A QoS 1 redelivery becomes a duplicate Kafka record. Include the original MQTT message identifier in the record so downstream can dedupe, and make the downstream sink idempotent anyway.

Exactly-once. Not available across this boundary. Kafka transactions cover Kafka to Kafka. Getting exactly-once semantics from an MQTT publish to a Kafka record requires a distributed commit between the broker’s session state and the log’s transaction state, which is not something either protocol offers you.

Choose at-least-once with an idempotent sink, or accept duplicates knowingly. Do not spend a quarter trying to build a two-phase commit across a device protocol boundary.

Reconnect storms

This is the failure that takes down ingestion platforms, and it is a handshake problem.

Five million devices with an exponential backoff and no jitter will reconnect in synchronised waves. Add a broker node restart or a mobile network outage, and a million devices can attempt to reconnect within the same few seconds.

1,000,000 reconnects in 10 seconds = 100,000 TLS handshakes/sec

A TLS 1.3 handshake is CPU-bound on the server side, and the achievable rate is on the order of a few thousand per core depending on hardware. One hundred thousand per second is therefore tens of cores consumed entirely by handshakes, on a node that was sized for message throughput.

Three mitigations, and all three are needed.

Jitter. Full jitter on the backoff, not exponential backoff with no randomness. Exponentially growing intervals still leave the herd synchronised, because the intervals are identical across devices.

Staged reconnection limits per broker node. A connection rate limiter, and a node that sheds excess connections rather than accepting them and timing out. Rejected devices retry; a storm of TLS handshakes is not.

Spread the fleet’s schedule. If devices report every thirty seconds, jitter the phase so they are not all reporting on the same second boundary in the first place. This is a firmware design decision and it is the cheapest fix available.

Watch connection churn per second as a first-class metric. A stable fleet has a low churn rate, and a rising churn rate is a precursor to a storm rather than a symptom of one.

Device commands, the other direction

Telemetry is one direction and easy. Commands are the reverse, and they are where a data platform becomes an operations platform.

A command needs a correlation identifier so the device can report the result, and it needs an expiry.

send unlock
   -> arrives three hours later on a device that has been sitting in a depot
   -> opens a door nobody is standing at

That is a real incident, and the fix is a timestamp and a TTL on every command, enforced at the device. An expired command is dropped and reported as expired rather than executed.

Commands also need a durable record before they are sent, not after, so that the platform can answer “did this get delivered, and what happened” when someone asks three months later. That record belongs in a transactional store, written with an outbox, and published by a separate dispatcher.

And commands must not be fire-and-forget over a QoS 0 channel. A command that is not acknowledged is not a command that is pending, it is a command whose fate is unknown, and the operator interface should show those differently from both succeeded and failed.

What ten million per second actually requires

The design above holds at 166,000 per second and has headroom to 1.7 million at peak. If the real requirement is ten million messages per second, several things stop being tunable.

The fleet is bigger. Ten million per second at one message per thirty seconds is three hundred million devices. That is not a broker sizing question, it is a business question, and it changes the provisioning and cost model completely.

The bridge is now a stream processing problem. At ten million per second, a single bridge process cannot batch fast enough, and you need many bridge instances consuming partitioned topic ranges, with the ordering key preserved.

Message size becomes the lever. At 200 bytes, ten million per second is 2 GB per second, which is 172 TB per day. Dropping to 30 bytes through binary encoding takes that to 26 TB per day, and it is the single biggest available saving.

You will need regional brokers. A hundred broker nodes is a manageable fleet. A thousand is an organisational problem, and users in different continents should not have their messages crossing an ocean to reach a consumer.

Exactly-once stops being a discussion. At this volume, the engineering effort goes to at-least-once plus idempotency, permanently.

The structural decisions survive the jump. Connections, session state, provisioning, key choice, and reconnect behaviour are the same problems at ten million per second. Only the constants move.

Failure stories worth testing

Restart one broker node holding 50,000 devices

Watch reconnect rate, handshake CPU, and how long it takes to re-establish. This is the reconnect storm in miniature, and it is the test that tells you whether your backoff has jitter in it.

Take the Kafka cluster away

Decide what the bridge does. Spill to disk with a bounded budget, or drop with a declared policy. Verify the broker’s own queues do not grow past the point where it starts evicting sessions, which looks exactly like a broker bug and is not.

Reboot ten thousand devices simultaneously

Force the condition rather than waiting for it. Firmware that reconnects immediately on boot will synchronise, and the fix is a randomised boot delay.

Expire a certificate across a subset of devices

Confirm the broker rejects them cleanly rather than dropping the connection in a way that causes an immediate reconnect loop. A rejected device that retries without backoff is a self-inflicted storm.

Publish to another device’s topic

Confirm the bridge rejects it. If it does not, the entire device identity model is decorative.

Fill a device’s session queue

Set a long session expiry on QoS 1 and take a broker offline for a day. Measure how much memory the queued messages actually consumed.

Duplicate every message

Replay the same MQTT message identifiers and confirm the downstream sink is unaffected.

A production-ready architecture

  +--------------+        MQTT over TLS 1.3        +------------------+
  | 5M devices   | ------------------------------> | Broker cluster   |
  | QoS 0        |  one-way, no session persistence | 100 nodes        |
  | clean session|                                  | 50k conns/node   |
  +--------------+                                  +--------+---------+
                                                           |
                              per-tenant ACLs, connection identity,
                              client-side topic prefix
                                                           |
                        +----------------------------------+-------------------+
                        |                                                      |
             +----------v-----------+                          +--------------v---------+
             | Bridge fleet        |                          | Command dispatcher    |
             | batch, key, dedupe  |                          | durable record first  |
             +----------+-----------+                          +------------+----------+
                        |                                                   |
          +-------------+-------------+                        +------------v---------+
          |             |             |                        | Command store       |
   +------v-----+ +-----v------+ +----v---------+              | (outbox pattern)    |
   | Kafka      | | Time series| | Object store |              +----------------------+
   | telemetry  | | hot tier   | | raw archive  |
   +------------+ +------------+ +--------------+

Control plane, deliberately separate from the data plane:

  PKI / certificate authority   -> short-lived certs, no revocation list
  Device registry                -> tenant, hardware root, firmware version
  Tenant policy                  -> topic prefixes, QoS rules, rate limits

A sensible delivery checklist:

  1. Size the broker tier by concurrent connections and connection churn, not by messages per second.
  2. Fix the message key at the bridge, because MQTT ordering stops at the connection.
  3. Choose QoS per data class, and default high-frequency telemetry to QoS 0.
  4. Bound session expiry so a long outage cannot become terabytes of queue.
  5. Enforce tenant boundaries at the broker and treat the topic’s device id as untrusted.
  6. Validate the publishing identity against the connection, not against the topic.
  7. Add jitter to every backoff, and rate-limit reconnection per node.
  8. Decide the bridge’s behaviour when Kafka is unavailable, and make the decision visible.
  9. Put a durable record and an expiry on every command, and enforce the expiry on the device.
  10. Alert on active connections, churn, and session memory, not on broker CPU.

Common mistakes

Mistake What actually happens Better decision
Sizing brokers by message rate Connections hit the ceiling long before throughput does Size by concurrent connections and per-connection memory
QoS 1 with persistent sessions everywhere Buffering one offline day is 2.88 TB of queue QoS 0 for telemetry, QoS 1 for critical events only
Trusting QoS 2 for exactly-once The guarantee ends at the broker and says nothing about Kafka Choose at-least-once with an idempotent sink
Relying on MQTT ordering Two devices on one topic have no ordering relationship Key Kafka records by device and metric
Trusting the device id in the topic Any device can publish as any other device Derive identity from the connection
A shared fleet credential One leak compromises everything and cannot be revoked Per-device certificate from a hardware root
No jitter in backoff A million devices reconnect in the same second Full jitter on every retry
Synchronised reporting phases A firmware choice manufactures a permanent storm Randomise the report phase at the device
Unbounded bridge buffering The broker evicts sessions and the symptom looks like a broker bug Bound the spill budget and declare the drop policy
Per-device topic ACLs Millions of pattern rules evaluated on every publish Enforce the prefix on the client, the tenant at the broker
Fleet-wide wildcard subscriptions Broker memory and per-publish cost scale with matches Use a stream processor instead of a subscription
Commands with no expiry A three-hour-old unlock executes when the device reconnects Timestamp and TTL every command, enforced on the device
Long-lived certificates with no revocation story A compromised cert is valid for months Short-lived certs from a hardware root, no CRL
Monitoring broker CPU only Churn rises for hours before the storm lands Alert on connections, churn, and session memory

The complete story in one minute

A device connects over MQTT with mutual TLS, presenting a certificate it fetched from a short-lived authority using a hardware root. It publishes telemetry at QoS 0 with a clean session, so nothing is buffered on the broker and nothing needs replaying, because a missing reading will be superseded thirty seconds later.

The broker cluster holds fifty thousand connections per node across a hundred nodes, enforces tenant boundaries from the connection identity, and does not attempt per-device topic rules. The device constructs only its own topic prefix.

A bridge fleet consumes from the broker, batches messages over a short window, and produces to Kafka keyed by device and metric, which is where per-device ordering actually comes from. Each record carries the original MQTT message identifier so the sink can deduplicate, and the sink is idempotent regardless, because the bridge is at-least-once and there is no exactly-once option across this boundary.

Commands travel the other way through a dispatcher that writes a durable record before sending, attaches a correlation identifier and an expiry, and treats an unacknowledged command as a third state rather than as either success or failure.

A reconnect storm is prevented rather than survived, through full jitter on backoff, per-node connection rate limits, and a firmware report phase that is randomised at the device.

That is the whole path:

device -> TLS -> broker -> bridge -> keyed Kafka records -> idempotent sinks
                     |
                     +-> connection state, session expiry, reconnect control

The hard parts were never the messages per second. They were five million connections with bounded state, a key that restores ordering the protocol threw away, and the discipline to default to the cheapest delivery guarantee that the data actually needs.

What this team still owns

MQTT QoS describes delivery between one sender and one receiver; it does not give end-to-end exactly-once processing through the broker, stream, database, and alert. Carry a stable device ID, boot/session ID, monotonic sequence where the device can provide one, schema version, device timestamp, and receipt timestamp. Deduplicate inside an explicit window, treat sequence reset as a first-class event, and publish the loss/duplication contract per telemetry class instead of assigning QoS 2 to everything.

Technical references

Keep reading
Browse everything