Consuming MQTT with aiomqtt¶
aiomqtt wraps the paho MQTT client in an asyncio interface: connect with async with, subscribe, and iterate client.messages. That simplicity hides three decisions that determine whether messages are lost: how fast messages are handled, what session the client keeps while disconnected, and what happens when the connection drops. Measured on Python 3.14 with aiomqtt 2.5.1 against Mosquitto 2.1.2 in Docker: a handler taking 10 ms per message, receiving 2,000 messages per second, handled 567 in 6 seconds while 9,432 waited in the client's unbounded queue; capping the queue at 500 discarded the excess instead, and 20 concurrent handlers kept up with all 10,000. A subscriber that disconnected while 100 QoS 1 messages were published received none of them on reconnecting with a clean session, and all 100 with clean_session=False. Throughput fell with the QoS level: 10,086 messages per second end to end at QoS 0, 3,994 at QoS 1 and 2,904 at QoS 2. And when the broker restarted, iteration raised MqttError: Disconnected during message iteration; a reconnect loop was back 1.63 s after the restart began. This guide builds a consumer that handles all three.
Prerequisites¶
- aiomqtt and an MQTT broker; the examples use Mosquitto.
- Poison-message handling, from handling poison messages in async consumers.
- The topic overview, Message Brokers & Event Streams.
1. Choose the QoS level by what loss costs¶
MQTT's quality-of-service level sets the delivery guarantee per subscription and per publish: 0 is at most once, 1 at least once, 2 exactly once between client and broker. Each step adds protocol round trips:
async with aiomqtt.Client("broker", 1883, identifier="meter-ingest") as client:
await client.subscribe("meters/+/reading", qos=1)
async for message in client.messages:
await handle(message)
Measured with 20,000 messages of 100 bytes from one publisher to one subscriber: QoS 0 ran at 10,299 publishes per second and 10,086 end to end; QoS 1 at 3,994; QoS 2 at 2,904. All 20,000 arrived at every level on a healthy local connection; the levels differ in what happens when something fails. QoS 1 with idempotent handlers — duplicates possible, loss not — is the usual choice; QoS 2 costs another 27% for a guarantee that ends at the broker, not at your database.
Verify: the QoS of each subscription is chosen deliberately, and handlers for QoS 1 tolerate duplicates.
2. Keep up with the message rate¶
client.messages is fed by aiomqtt from the network as fast as messages arrive, regardless of how fast the async for body runs. Measured with a handler that awaited 10 ms per message — 100 per second — and a publisher sending 2,000 per second for 5 seconds:
async for message in client.messages:
await asyncio.sleep(0.01) # stand-in for a 10 ms handler
After 6 seconds, 567 messages had been handled and 9,432 were waiting in the client's queue, which is unbounded by default. In a long-running service that queue grows until the process runs out of memory, or the backlog makes every message minutes old. Setting max_queued_incoming_messages=500 bounded the queue at 500, and aiomqtt discarded the rest with a log line, Message queue is full. Discarding message. — bounded memory, but silent loss for QoS 1 messages the broker considers delivered.
The fix is to handle messages concurrently. Several tasks can iterate the same client.messages:
async def worker(client):
async for message in client.messages:
await handle(message)
async with aiomqtt.Client("broker", 1883, identifier="meter-ingest") as client:
await client.subscribe("meters/+/reading", qos=1)
async with asyncio.TaskGroup() as tg:
for _ in range(20):
tg.create_task(worker(client))
Measured: 20 workers handled all 10,000 messages with a maximum backlog of 730; 40 workers kept it at 26. Size the worker count as message rate times handling time, with headroom, as in rate-limited worker pools.
Verify: under peak message rate, the length of client.messages stays bounded and near zero.
3. Keep a session across disconnects¶
What a subscriber receives after a disconnect depends on its session. With a clean session, the broker forgets the subscription when the client disconnects; with a persistent session, it keeps the subscription and queues QoS 1 and 2 messages for the client's identifier:
client = aiomqtt.Client(
"broker", 1883,
identifier="meter-ingest-1", # stable: the session is keyed by it
clean_session=False, # MQTT 3.1.1; use clean_start with MQTT 5
)
Measured by subscribing, disconnecting, publishing messages 100–199, and reconnecting with the same identifier without subscribing again: with a clean session at QoS 1, 0 of the 100 arrived. With clean_session=False at QoS 0, 0 of 100 — QoS 0 messages are not queued for offline clients. With clean_session=False at QoS 1, all 100 arrived on reconnect. The identifier must be stable across restarts and unique per consumer instance; two clients with the same identifier disconnect each other. The broker's queue for an offline client is bounded by its own settings — max_queued_messages in Mosquitto — so a long outage can still lose messages.
Verify: a test that stops the consumer, publishes, and restarts it receives every QoS 1 message published while it was down.
4. Reconnect in a loop¶
aiomqtt does not reconnect by itself: when the connection drops, iteration raises aiomqtt.MqttError, and the async with block ends. Wrap the client in a loop:
async def consume_forever(handle, interval: float = 1.0):
while True:
try:
async with aiomqtt.Client("broker", 1883, identifier="meter-ingest-1",
clean_session=False) as client:
await client.subscribe("meters/+/reading", qos=1)
async for message in client.messages:
await handle(message)
except aiomqtt.MqttError as e:
log.warning("mqtt disconnected: %s; reconnecting in %.1fs", e, interval)
await asyncio.sleep(interval)
Measured with docker restart -t 0 on the broker: the restart took 0.63 s, the subscriber saw one MqttError: Disconnected during message iteration, and it was connected and subscribed again 1.63 s after the restart began — the restart plus the 1-second interval. With a persistent session, messages published in that gap were delivered after reconnecting, as in step 3. Add jittered backoff for production, so a fleet of consumers does not reconnect in lockstep after a broker restart.
Verify: restarting the broker during consumption produces a reconnect within the expected interval and no lost QoS 1 messages.
5. Make handling idempotent¶
QoS 1 delivers at least once. A message whose handling finished but whose acknowledgement did not reach the broker before a disconnect is delivered again on reconnect. Make handlers safe to repeat:
async def handle(message: aiomqtt.Message):
reading = json.loads(message.payload)
await db.execute(
"INSERT INTO readings (meter_id, ts, value) VALUES ($1, $2, $3) "
"ON CONFLICT (meter_id, ts) DO NOTHING",
reading["meter"], reading["ts"], reading["value"],
)
A natural key from the message — meter and timestamp here — turns a duplicate into a no-op. Where the payload has no natural key, publishers can include an identifier, and the consumer can record processed identifiers with an expiry. A handler that raises on a malformed message should log and skip it, not crash the loop, or one bad message will reconnect the client forever; see handling poison messages in async consumers.
Verify: replaying a batch of messages through the handler leaves the stored state unchanged.
Verification¶
An aiomqtt consumer is reliable when:
- Handlers keep up: concurrent workers keep
client.messagesnear empty at peak rate. - Sessions persist: a stable
identifierwithclean_session=Falseand QoS 1. - Reconnection is explicit: a loop around the client catches
MqttErrorand backs off. - Handling is idempotent, so redeliveries are harmless.
Diagnostic Hook: when an MQTT consumer's memory grows steadily while the broker shows no backlog, measure len(client.messages). The messages have been delivered to the client and are queued in it — 9,432 after 6 seconds with a 10 ms handler at 2,000 messages per second.
Pitfalls & edge cases¶
- A slow handler in a single
async for. Measured: 9,432 messages queued in 6 s. - Capping the client queue as the fix. Messages were discarded instead.
- Clean sessions, or QoS 0, for data that matters. Measured: 0 of 100 offline messages delivered.
- Expecting automatic reconnection. aiomqtt raises
MqttError; the loop is yours.
Frequently Asked Questions¶
Does aiomqtt reconnect automatically?
No. Iteration raises aiomqtt.MqttError when the connection drops. Wrap the client in a while loop that catches MqttError and sleeps; it reconnected 1.63 s after a broker restart.
How do I receive MQTT messages sent while my client was offline?
Use a stable identifier, clean_session=False and QoS 1. That received all 100 offline messages; a clean session or QoS 0 received none.
Why is my aiomqtt consumer using more and more memory?
The handler is slower than the message rate and messages queue in the client. A 10 ms handler at 2,000 msg/s left 9,432 queued; 20 concurrent workers kept up.
How much slower is MQTT QoS 1 than QoS 0?
Here 3,994 against 10,086 messages per second end to end, and QoS 2 2,904, with one publisher and one subscriber on a local broker.
Related¶
- Message Brokers & Event Streams — up to the topic overview.
- Publishing with confirms in aio-pika — the publishing side of delivery guarantees.
- Network I/O & Protocol Handling — the section overview.