Skip to content

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

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.

aiomqtt 2.5.1 with Mosquitto 2.1.2 A grid of 6 rows by 2 columns. aiomqtt 2.5.1 with Mosquitto 2.1.2 scenario result QoS 0 / 1 / 2, 20,000 messages 10,086 / 3,994 / 2,904 msg/s, none lost 10 ms handler, 2,000 msg/s, one consumer 567 handled, 9,432 queued in the client same, max_queued_incoming_messages=500 565 handled, the rest discarded same, 20 concurrent handlers 10,000 handled, max backlog 730 offline during 100 QoS 1 messages, clean session 0 of 100 received offline, clean_session=False, QoS 1 100 of 100 received Python 3.14, local Docker broker.

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.

Messages handled out of 10,000 in 6 seconds 4 horizontal bars comparing 1 consumer (9,432 queued) with the others. Messages handled out of 10,000 in 6 seconds 1 consumer (9,432 queued) 567 1 consumer, queue capped at 500 (rest discarded) 565 20 workers (max backlog 730) 10,000 40 workers (max backlog 26) 10,000 The client queue absorbs what the handler cannot.

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.

A reliable aiomqtt consumer A flow of 5 stages. A reliable aiomqtt consumer Connect stable identifier, clean_session=False Subscribe QoS 1 Workers concurrent handlers, queue near zero Handle idempotent, duplicates possible MqttError back off, reconnect, get queued messages Session, concurrency and reconnection each close a different loss path.

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.messages near empty at peak rate.
  • Sessions persist: a stable identifier with clean_session=False and QoS 1.
  • Reconnection is explicit: a loop around the client catches MqttError and 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.