This release is a pre-release and may not be stable for production use.
MQTTium
MQTTium is a reliable, async-native MQTT client for modern Python.
It combines a synchronous protocol engine with an asyncio API, bounded
backpressure, durable inflight persistence, explicit delivery receipts, and
complete MQTT QoS state machines.
Status: beta (
0.2.0b3). The native Stable API tier follows the documented compatibility policy; Provisional extension surfaces may still evolve before the first stable release.
Features
- MQTT 3.1.1 and MQTT 5;
- QoS 0, 1 and 2 with one authoritative protocol state machine;
- TCP, TLS, WebSocket and Unix transports;
- reconnect and session replay;
- bounded callback and async-iterator delivery;
- immutable runtime statistics for queues, budgets, receipts and transports;
- manual acknowledgement;
- in-memory and SQLite inflight persistence;
- aggregate
publish_many()with bounded memory and measured throughput gains; - synchronous loop-bound
publish_nowait()for non-suspending native producers; - an additive Paho VERSION2 compatibility façade, isolated from the native API;
- inline type information for type checkers.
Python 3.11–3.14 is supported. MQTTium is licensed under Apache-2.0.
API stability
The native client is imported from mqttium.api; protocol enums and operational
exceptions are imported from mqttium. Advanced protocol, persistence,
transport, packet and Paho surfaces are classified separately as Provisional.
See docs/API-STABILITY.md for the exact Stable,
Provisional and Internal boundaries.
Installation
Install the current beta from PyPI:
python -m pip install --pre mqttium
To work on the unreleased development version:
git clone https://github.com/yoch/mqttium.git
cd mqttium
python -m pip install -e ".[dev,fuzz,release]"
Quick start
import asyncio
from mqttium.api import AsyncClient
async def main() -> None:
client = AsyncClient("demo-client")
await client.connect("127.0.0.1", 1883)
await client.subscribe("demo/#")
receipt = await client.publish("demo/hello", b"world", qos=1)
await receipt.wait()
await client.disconnect()
asyncio.run(main())
Non-suspending publishing
A producer already executing on the client's event loop can submit without creating or awaiting a coroutine:
receipt = client.publish_nowait("telemetry/device-1", payload, qos=1)
# Continue synchronously, then observe completion later if needed.
await receipt.wait()
publish_nowait() either admits the publication immediately or raises
FlowControlError; it never waits for engine or writer capacity. It returns the
normal PublishReceipt, so QoS 1/2 completion is observed in the same way as
with publish().
The method follows the same ownership rule as asyncio.Queue.put_nowait(): it
is intended for the client's owning event-loop thread, not as a generic
thread-safe API. Cross-thread synchronous callers should use an adapter such as
mqttium.compat.paho.Client, which coalesces submissions before handing a
bounded batch to the loop.
Batched publishing
publish_many() consumes iterables in bounded chunks and returns one aggregate
receipt rather than creating one task or event per message:
from mqttium.api import PublishMessage
batch = await client.publish_many(
PublishMessage("telemetry/device-1", payload, qos=1)
for payload in payloads
)
await batch.wait()
The retained paired A/B benchmark measured publisher-throughput geomean
improvements of 36.9% for QoS 0, 15.9% for QoS 1, and 7.0% for QoS 2
against equivalent individual publishing pipelines on the validated source
tree. Benchmark methodology and limitations are documented in
docs/BENCHMARKING.md.
Bounded memory
Every queue that can grow with application load is bounded by default, so a producer that outruns its broker is slowed down rather than allowed to exhaust the process:
client = AsyncClient(
max_pending_outbound_messages=10_000, # unfinished QoS 1/2 publications
max_pending_outbound_bytes=64 * 1024**2, # their logical topic+payload+properties
max_pending_inbound_bytes=64 * 1024**2, # persisted inbound QoS handshakes
max_pending_delivery_bytes=64 * 1024**2, # inbound messages awaiting a consumer
max_ingress_batch_bytes=1 * 1024**2, # decoded work before delivery is drained
publish_backpressure="wait", # or "error" to refuse immediately
)
publish() waits for capacity by default and raises FlowControlError under
publish_backpressure="error" or with nowait=True. A refusal is atomic: no
packet identifier is allocated and no store record is written. Pass None for
any limit to restore unbounded queueing.
The two inbound byte limits cover different lifetimes:
max_pending_inbound_bytes bounds QoS 2 and manually acknowledged QoS 1 data
retained by the session store, while max_pending_delivery_bytes bounds data
waiting for callback or iterator consumers. MQTT 5 reports a persisted-session
quota violation with reason 0x97; MQTT 3.1.1 closes the connection.
The writer also has an independent max_outbound_bytes=1 MiB default.
publish_nowait() refuses immediately when that wire queue is full, so callers
publishing large payloads should size this byte limit explicitly rather than
assuming max_outbound_messages is the only binding constraint.
Native QoS 0 uses its direct writer fast path only while on_publish is None.
Installing that callback deliberately routes publications through the standard
engine and EffectPump so completion callbacks retain their ordering semantics.
These defaults are new in 0.1.0a2; before them a QoS 1/2 producer could queue
until the 65 535 packet-identifier space was exhausted. See
docs/MIGRATION.md.
Runtime statistics
stats() returns a frozen snapshot without starting a sampler or emitting logs:
snapshot = client.stats()
print(snapshot.outbound.pending_bytes)
print(snapshot.inbound.inflight)
print(snapshot.writer.queued_bytes)
print(snapshot.delivery.pending_bytes)
Each section is produced by the component that owns the state — the two protocol
sessions, the effect and write pumps, the transport — and stats() only
assembles them. The snapshot also includes lifetime high-water marks, batching
decision counters, task state, receipt counts, decoder buffering and
WebSocket/stream transport buffers. It is intended to be called on the client's
owning event loop. See docs/API-STABILITY.md.
Validation
The release gates include:
- more than 500 unit tests;
- Mosquitto integration tests on Python 3.11, 3.12, 3.13 and 3.14;
- deterministic and Hypothesis-based fuzzing;
- Ruff formatting and linting;
- mypy validation and a PEP 561
py.typedmarker; - an 80% coverage gate;
- wheel, source-distribution and isolated-install validation;
- delivery, persistence, TCP, TLS and WAN-profile benchmarks.
A separate finalisation workflow runs short reconnect/backpressure soaks on
Linux for relevant pull requests. Extended Linux/macOS soaks and interoperability
campaigns against multiple brokers are available by manual dispatch. Their
acceptance criteria are documented in docs/STABILITY.md.
Documentation
docs/DESIGN.md— architecture and invariantsdocs/API-STABILITY.md— public API policy and deprecationsdocs/STABILITY.md— soak and interoperability campaigndocs/RELEASING.md— rehearsal, publication and failure handlingdocs/IMPLEMENTATION-GUIDE.md— protocol contractsdocs/COMPAT.md— Paho compatibility surfacedocs/MIGRATION.md— migration guidancedocs/BENCHMARKING.md— benchmark validity contractdocs/FUZZING.md— fuzzing strategydocs/ROADMAP.md— remaining stable-release workPROVENANCE.md— source history and licensing review
Contributing and security
See CONTRIBUTING.md. Report vulnerabilities through the
private process described in SECURITY.md.
Release files for mqttium 0.2.0b3
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| mqttium-0.2.0b3.tar.gz | 290.6 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| mqttium-0.2.0b3-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 413.3 kB
Release files / mqttium-0.2.0b3.tar.gz
| Download URL | mqttium-0.2.0b3.tar.gz |
|---|---|
| Size | 290.6 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
4ac18e3e95fda481ef897aa2c9d6a4d9f69720c3e009bbd1c187ad6415944e05
|
|
BLAKE2b-256 checksum How to use checksums |
a1623fd84fa2f302040b8953a1b9385875a147d6669c3264c304f393671d487a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Aug 9, 2026.
Transparency logRelease files / mqttium-0.2.0b3-py3-none-any.whl
| Download URL | mqttium-0.2.0b3-py3-none-any.whl |
|---|---|
| Size | 122.7 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
7c282d5ca9d01c34ec6f3a6e52bc5f1c9cd647ace86a961673df8afe867eb91f
|
|
BLAKE2b-256 checksum How to use checksums |
d21eabc37eaa77c13141f5262349994021a9f4c464506b24e3a9b72e8a0ba10f
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Aug 9, 2026.
Transparency log