MQTT
A lightweight publish-subscribe messaging protocol designed for constrained devices and low-bandwidth, high-latency networks.
MQTT Cheatsheet
A lightweight publish-subscribe messaging protocol designed for constrained devices and low-bandwidth, high-latency networks.
Overview
MQTT (Message Queuing Telemetry Transport) is an OASIS standard messaging protocol built on TCP/IP, optimised for IoT and M2M communication. It uses a publish-subscribe pattern where clients connect to a central broker that handles message routing between publishers and subscribers.
graph TB
subgraph "MQTT Architecture"
Broker[MQTT Broker]
subgraph "Publishers"
P1[Sensor 1]
P2[Sensor 2]
P3[Device]
end
subgraph "Subscribers"
S1[Dashboard]
S2[Database]
S3[Alert Service]
end
P1 -->|Publish| Broker
P2 -->|Publish| Broker
P3 -->|Publish| Broker
Broker -->|Subscribe| S1
Broker -->|Subscribe| S2
Broker -->|Subscribe| S3
end
style Broker fill:#e1f5fe
Broker Setup and Configuration
Key Concepts
- Broker: Central server that receives messages from publishers and routes them to subscribers
- Listener: Network interface and port combination for accepting connections
- ACL (Access Control List): Rules defining which clients can publish/subscribe to which topics
- Persistence: Storage of messages and session state for QoS 1/2 and retained messages
- Bridge: Connection between two MQTT brokers for federation
Common Patterns
Mosquitto Installation and Basic Configuration:
# Install Mosquitto broker (Ubuntu/Debian)
apt-get update && apt-get install -y mosquitto mosquitto-clients
# Install on macOS
brew install mosquitto
# Install on Docker
docker run -d --name mosquitto \
-p 1883:1883 \
-p 9001:9001 \
-v /path/to/mosquitto.conf:/mosquitto/config/mosquitto.conf \
-v /path/to/data:/mosquitto/data \
-v /path/to/log:/mosquitto/log \
eclipse-mosquitto:latest
Basic Mosquitto Configuration (mosquitto.conf):
# Network listeners
listener 1883
listener 9001
protocol websockets
# Persistence
persistence true
persistence_location /var/lib/mosquitto/
# Logging
log_dest file /var/log/mosquitto/mosquitto.log
log_type all
# Security - disable anonymous access
allow_anonymous false
password_file /etc/mosquitto/passwd
# Connection limits
max_connections 1000
max_inflight_messages 20
max_queued_messages 1000
# Message size limit (bytes)
message_size_limit 1048576
# Keepalive settings (max keepalive the broker will allow, in seconds)
max_keepalive 65535
# ACL file
acl_file /etc/mosquitto/acl
Authentication Setup:
# Create password file
mosquitto_passwd -c /etc/mosquitto/passwd username1
mosquitto_passwd -b /etc/mosquitto/passwd username2 password2
# Add user to existing file
mosquitto_passwd /etc/mosquitto/passwd newuser
ACL Configuration (/etc/mosquitto/acl):
# User-specific rules
user admin
topic readwrite #
user sensor1
topic write sensors/temperature/+
topic read commands/sensor1/#
user dashboard
topic read sensors/#
topic read status/#
# Pattern-based rules (using %u for username, %c for client ID)
pattern readwrite users/%u/#
pattern read $SYS/#
TLS/SSL Configuration:
# TLS listener
listener 8883
cafile /etc/mosquitto/certs/ca.crt
certfile /etc/mosquitto/certs/server.crt
keyfile /etc/mosquitto/certs/server.key
require_certificate false
# Require client certificates (mTLS)
listener 8884
cafile /etc/mosquitto/certs/ca.crt
certfile /etc/mosquitto/certs/server.crt
keyfile /etc/mosquitto/certs/server.key
require_certificate true
use_identity_as_username true
Examples
Docker Compose with Mosquitto:
version: '3.8'
services:
mosquitto:
image: eclipse-mosquitto:latest
container_name: mqtt-broker
ports:
- "1883:1883"
- "9001:9001"
volumes:
- ./mosquitto/config:/mosquitto/config
- ./mosquitto/data:/mosquitto/data
- ./mosquitto/log:/mosquitto/log
restart: unless-stopped
Bridge Configuration (connecting two brokers):
# Bridge to remote broker
connection bridge-to-cloud
address cloud.example.com:8883
topic sensors/# out 1
topic commands/# in 1
bridge_cafile /etc/mosquitto/certs/ca.crt
bridge_certfile /etc/mosquitto/certs/client.crt
bridge_keyfile /etc/mosquitto/certs/client.key
remote_username bridge-user
remote_password bridge-pass
start_type automatic
cleansession false
notifications true
Topic Structure and Wildcards
Key Concepts
graph TD
subgraph "Topic Hierarchy"
Root["/"]
Root --> Home["home"]
Root --> Office["office"]
Home --> Living["living-room"]
Home --> Bedroom["bedroom"]
Living --> Temp1["temperature"]
Living --> Humidity1["humidity"]
Bedroom --> Temp2["temperature"]
Bedroom --> Light["light"]
end
subgraph "Example Topics"
T1["home/living-room/temperature"]
T2["home/bedroom/light"]
T3["office/floor1/meeting-room/occupancy"]
end
- Topic: UTF-8 string used to filter messages (case-sensitive)
- Topic Level: Segments separated by forward slashes (
/) - Single-level Wildcard (
+): Matches exactly one topic level - Multi-level Wildcard (
#): Matches any number of levels (must be last character) - $SYS Topics: Reserved broker statistics topics (not matched by
#)
Common Patterns
Topic Naming Best Practices:
# Good topic structure
<domain>/<location>/<device-type>/<device-id>/<measurement>
# Examples
home/living-room/sensor/temp-01/temperature
factory/line-1/robot/arm-02/status
vehicle/fleet-a/truck-123/gps/position
# Avoid starting with /
WRONG: /sensors/temperature
RIGHT: sensors/temperature
# Use lowercase and hyphens
WRONG: Sensors/Temperature_Value
RIGHT: sensors/temperature-value
Wildcard Subscriptions:
# Single-level wildcard (+)
# Matches: home/living-room/temperature, home/bedroom/temperature
# Does NOT match: home/living-room/sensor/temperature
mosquitto_sub -t "home/+/temperature"
# Matches any device in kitchen
mosquitto_sub -t "home/kitchen/+/status"
# Multi-level wildcard (#)
# Matches ALL topics under home/
mosquitto_sub -t "home/#"
# Matches: sensors/temp/1, sensors/temp/1/value, sensors/humidity
mosquitto_sub -t "sensors/#"
# Combining wildcards
# Matches: home/room1/device1/status, home/room2/device2/status
mosquitto_sub -t "home/+/+/status"
# Subscribe to all topics (be careful with this!)
mosquitto_sub -t "#"
# Subscribe to broker statistics
mosquitto_sub -t "\$SYS/#"
Examples
Python Topic Handling:
paho-mqtt 2.x note: the Python examples below use the 1.x callback style. Since paho-mqtt 2.0 (Feb 2024), a bare
mqtt.Client()raisesValueError— you must pass an API version, e.g.mqtt.Client(mqtt.CallbackAPIVersion.VERSION2), and theon_connect/on_disconnectcallbacks gainreason_code/propertiesparameters. On 1.x the calls below work as written; on 2.x add theCallbackAPIVersionargument and update the callback signatures.
import paho.mqtt.client as mqtt
import re
def on_message(client, userdata, msg):
# Parse topic structure
parts = msg.topic.split('/')
if len(parts) >= 4:
location = parts[0]
room = parts[1]
device = parts[2]
measurement = parts[3]
print(f"Location: {location}, Room: {room}")
print(f"Device: {device}, Measurement: {measurement}")
print(f"Value: {msg.payload.decode()}")
client = mqtt.Client()
client.on_message = on_message
client.connect("localhost", 1883)
# Subscribe to multiple patterns
client.subscribe([
("home/+/temperature", 1),
("home/+/humidity", 1),
("alerts/#", 2)
])
client.loop_forever()
Quality of Service (QoS) Levels
Key Concepts
sequenceDiagram
participant Publisher
participant Broker
participant Subscriber
Note over Publisher,Subscriber: QoS 0 - At Most Once
Publisher->>Broker: PUBLISH
Broker->>Subscriber: PUBLISH
Note over Publisher,Subscriber: QoS 1 - At Least Once
Publisher->>Broker: PUBLISH
Broker->>Subscriber: PUBLISH
Subscriber->>Broker: PUBACK
Broker->>Publisher: PUBACK
Note over Publisher,Subscriber: QoS 2 - Exactly Once
Publisher->>Broker: PUBLISH
Broker->>Publisher: PUBREC
Publisher->>Broker: PUBREL
Broker->>Subscriber: PUBLISH
Subscriber->>Broker: PUBACK
Broker->>Publisher: PUBCOMP
| QoS Level | Name | Delivery Guarantee | Use Case |
|---|---|---|---|
| 0 | At most once | Fire and forget, may be lost | Sensor readings where occasional loss is acceptable |
| 1 | At least once | Guaranteed delivery, may duplicate | Commands, alerts where duplicates are tolerable |
| 2 | Exactly once | Guaranteed single delivery | Financial transactions, critical operations |
- Downgrade: Broker delivers at the minimum of publisher QoS and subscriber QoS
- Inflight Messages: QoS 1/2 messages awaiting acknowledgement
- Message ID: 16-bit identifier for QoS 1/2 message tracking
Common Patterns
CLI Examples with QoS:
# Publish with QoS 0 (default)
mosquitto_pub -t "sensors/temp" -m "22.5"
# Publish with QoS 1
mosquitto_pub -t "alerts/critical" -m "Temperature exceeded" -q 1
# Publish with QoS 2
mosquitto_pub -t "transactions/payment" -m '{"id":"123","amount":99.99}' -q 2
# Subscribe with QoS 1
mosquitto_sub -t "sensors/#" -q 1
# Subscribe with QoS 2
mosquitto_sub -t "commands/#" -q 2
Python QoS Handling:
import paho.mqtt.client as mqtt
def on_publish(client, userdata, mid):
print(f"Message {mid} published")
def on_message(client, userdata, msg):
print(f"Received (QoS {msg.qos}): {msg.topic} = {msg.payload.decode()}")
client = mqtt.Client()
client.on_publish = on_publish
client.on_message = on_message
client.connect("localhost", 1883)
# Publish with different QoS levels
# QoS 0 - Fire and forget
client.publish("sensors/temp", "22.5", qos=0)
# QoS 1 - At least once
result = client.publish("alerts/warning", "High temperature", qos=1)
result.wait_for_publish() # Block until published
# QoS 2 - Exactly once
info = client.publish("commands/critical", "shutdown", qos=2)
info.wait_for_publish()
# Subscribe with QoS
client.subscribe("sensors/#", qos=1)
client.loop_forever()
Examples
QoS Level Selection Guide:
# QoS 0: Non-critical, frequent data
# - Temperature readings every second
# - GPS coordinates from vehicles
# - Log messages
client.publish("telemetry/gps", position_data, qos=0)
# QoS 1: Important data where duplicates are acceptable
# - Device status updates
# - Alerts and notifications
# - Configuration changes
client.publish("devices/status", json.dumps(status), qos=1)
# QoS 2: Critical operations requiring exactly-once delivery
# - Payment transactions
# - Machine control commands
# - State changes in distributed systems
client.publish("control/valve/open", "1", qos=2)
Retained Messages
Key Concepts
sequenceDiagram
participant Publisher
participant Broker
participant NewSubscriber
Publisher->>Broker: PUBLISH (retained=true)
Note over Broker: Stores last retained message
Note over NewSubscriber: Subscribes later
NewSubscriber->>Broker: SUBSCRIBE
Broker->>NewSubscriber: PUBLISH (retained message)
Note over Publisher,NewSubscriber: New messages replace retained
Publisher->>Broker: PUBLISH (retained=true, new value)
Note over Broker: Updates stored message
- Retained Message: Last message stored by broker, sent to new subscribers immediately
- One Per Topic: Only one retained message per topic (new replaces old)
- Clear Retained: Publish empty payload with retain flag to delete
- Retain Flag: Boolean indicating message should be stored
Common Patterns
CLI Examples:
# Publish retained message
mosquitto_pub -t "devices/thermostat/status" -m "online" -r
# Publish retained JSON
mosquitto_pub -t "config/device1" -m '{"interval":60,"threshold":25}' -r
# Clear retained message (empty payload)
mosquitto_pub -t "devices/thermostat/status" -m "" -r -n
# Subscribe and receive retained messages
mosquitto_sub -t "devices/#" -v
Python Retained Messages:
import paho.mqtt.client as mqtt
import json
client = mqtt.Client()
client.connect("localhost", 1883)
# Publish device status as retained
status = {
"online": True,
"firmware": "2.1.0",
"last_boot": "2024-01-15T10:30:00Z"
}
client.publish(
"devices/sensor-01/status",
json.dumps(status),
qos=1,
retain=True # Store as retained message
)
# Publish configuration as retained
config = {
"sample_rate": 1000,
"threshold_min": 10,
"threshold_max": 50
}
client.publish(
"config/sensor-01",
json.dumps(config),
qos=1,
retain=True
)
# Clear a retained message
client.publish("devices/old-device/status", "", retain=True)
client.disconnect()
Examples
Common Use Cases for Retained Messages:
# Device presence/status
# New dashboards immediately know device state
client.publish("devices/pump-01/status", "running", retain=True)
# Current configuration
# New services get latest config on subscribe
client.publish("config/global/settings", json.dumps(settings), retain=True)
# Last known sensor value
# Useful for infrequently updating sensors
client.publish("sensors/outdoor/temperature", "18.5", retain=True)
# System state
client.publish("system/mode", "production", retain=True)
Subscriber Handling Retained Messages:
def on_message(client, userdata, msg):
# Check if message is retained
if msg.retain:
print(f"Retained message: {msg.topic}")
# Handle initial state differently if needed
process_initial_state(msg.topic, msg.payload)
else:
print(f"Live message: {msg.topic}")
process_live_update(msg.topic, msg.payload)
Last Will and Testament (LWT)
Key Concepts
sequenceDiagram
participant Client
participant Broker
participant Subscribers
Client->>Broker: CONNECT (with LWT)
Note over Broker: Stores LWT message
Broker->>Client: CONNACK
Note over Client: Normal operation
Client->>Broker: PUBLISH messages
Note over Client: Unexpected disconnect
Client--xBroker: Connection lost
Note over Broker: Detects client gone
Broker->>Subscribers: PUBLISH (LWT message)
- LWT (Last Will and Testament): Message published by broker when client disconnects unexpectedly
- Set at Connect: Configured during CONNECT, cannot be changed after
- Trigger Conditions: Network failure, keepalive timeout, protocol error (not clean disconnect)
- Combined with Retain: LWT can be retained to persist offline status
Common Patterns
CLI with LWT:
# Connect with Last Will message
mosquitto_sub -t "test" \
--will-topic "clients/sensor-01/status" \
--will-payload "offline" \
--will-qos 1 \
--will-retain
# Another client subscribing to status
mosquitto_sub -t "clients/+/status" -v
Python LWT Configuration:
import paho.mqtt.client as mqtt
import json
def on_connect(client, userdata, flags, rc):
if rc == 0:
print("Connected successfully")
# Publish online status
client.publish(
"clients/sensor-01/status",
json.dumps({"status": "online", "timestamp": time.time()}),
qos=1,
retain=True
)
else:
print(f"Connection failed: {rc}")
def on_disconnect(client, userdata, rc):
if rc != 0:
print(f"Unexpected disconnect: {rc}")
client = mqtt.Client(client_id="sensor-01")
# Configure Last Will and Testament
lwt_payload = json.dumps({
"status": "offline",
"reason": "unexpected_disconnect"
})
client.will_set(
topic="clients/sensor-01/status",
payload=lwt_payload,
qos=1,
retain=True # Persist offline status
)
client.on_connect = on_connect
client.on_disconnect = on_disconnect
client.connect("localhost", 1883, keepalive=60)
client.loop_forever()
Examples
Device Presence Pattern:
import paho.mqtt.client as mqtt
import json
import time
class MQTTDevice:
def __init__(self, device_id, broker_host):
self.device_id = device_id
self.client = mqtt.Client(client_id=device_id)
# Status topic
self.status_topic = f"devices/{device_id}/status"
# Configure LWT for offline detection
self.client.will_set(
self.status_topic,
json.dumps({"online": False, "last_seen": time.time()}),
qos=1,
retain=True
)
self.client.on_connect = self._on_connect
self.client.connect(broker_host, 1883, keepalive=30)
def _on_connect(self, client, userdata, flags, rc):
# Announce online status
client.publish(
self.status_topic,
json.dumps({"online": True, "connected_at": time.time()}),
qos=1,
retain=True
)
def graceful_disconnect(self):
# Publish offline status before clean disconnect
self.client.publish(
self.status_topic,
json.dumps({"online": False, "reason": "shutdown"}),
qos=1,
retain=True
)
self.client.disconnect()
# Usage
device = MQTTDevice("sensor-01", "localhost")
device.client.loop_start()
# Do work...
# Clean shutdown
device.graceful_disconnect()
Fleet Monitoring with LWT:
def on_message(client, userdata, msg):
status = json.loads(msg.payload)
device_id = msg.topic.split('/')[1]
if status.get('online'):
print(f"Device {device_id} is ONLINE")
# Update dashboard, clear alerts
else:
print(f"Device {device_id} is OFFLINE")
# Trigger alert, update monitoring
# Subscribe to all device status
client.subscribe("devices/+/status", qos=1)
Client Libraries and Tools
Key Concepts
| Language | Library | Features |
|---|---|---|
| Python | paho-mqtt | Full featured, callbacks, threading |
| JavaScript | mqtt.js | Browser and Node.js support |
| Go | paho.mqtt.golang | Concurrent, connection management |
| Java | Eclipse Paho | Android support, extensive features |
| C/C++ | Mosquitto | Lightweight, embedded systems |
| Rust | rumqtt | Async/await, performance |
Common Patterns
Python (paho-mqtt):
import paho.mqtt.client as mqtt
import json
import ssl
# Callback functions
def on_connect(client, userdata, flags, rc):
codes = {
0: "Connected",
1: "Incorrect protocol version",
2: "Invalid client ID",
3: "Server unavailable",
4: "Bad credentials",
5: "Not authorised"
}
print(f"Connection result: {codes.get(rc, 'Unknown')}")
if rc == 0:
# Subscribe on connect to handle reconnection
client.subscribe([
("sensors/#", 1),
("commands/#", 2)
])
def on_message(client, userdata, msg):
print(f"{msg.topic}: {msg.payload.decode()}")
def on_disconnect(client, userdata, rc):
if rc != 0:
print(f"Unexpected disconnect, will auto-reconnect")
# Create client
client = mqtt.Client(
client_id="my-client",
clean_session=True,
protocol=mqtt.MQTTv311
)
# Set callbacks
client.on_connect = on_connect
client.on_message = on_message
client.on_disconnect = on_disconnect
# Authentication
client.username_pw_set("username", "password")
# TLS configuration
client.tls_set(
ca_certs="/path/to/ca.crt",
certfile="/path/to/client.crt",
keyfile="/path/to/client.key",
tls_version=ssl.PROTOCOL_TLSv1_2
)
# Connect with automatic reconnect
client.connect("broker.example.com", 8883, keepalive=60)
# Non-blocking loop
client.loop_start()
# Publish messages
client.publish("sensors/temp", "22.5", qos=1)
# Blocking loop (alternative)
# client.loop_forever()
JavaScript/Node.js (mqtt.js):
const mqtt = require('mqtt');
// Connection options
const options = {
clientId: 'my-client-' + Math.random().toString(16).substr(2, 8),
clean: true,
connectTimeout: 4000,
username: 'user',
password: 'pass',
reconnectPeriod: 1000,
// TLS options
// ca: fs.readFileSync('/path/to/ca.crt'),
// cert: fs.readFileSync('/path/to/client.crt'),
// key: fs.readFileSync('/path/to/client.key'),
// Last Will
will: {
topic: 'clients/my-client/status',
payload: JSON.stringify({ online: false }),
qos: 1,
retain: true
}
};
const client = mqtt.connect('mqtt://localhost:1883', options);
client.on('connect', () => {
console.log('Connected');
// Publish online status
client.publish('clients/my-client/status',
JSON.stringify({ online: true }),
{ qos: 1, retain: true }
);
// Subscribe to topics
client.subscribe({
'sensors/#': { qos: 1 },
'commands/#': { qos: 2 }
}, (err, granted) => {
if (err) console.error('Subscribe error:', err);
else console.log('Subscribed:', granted);
});
});
client.on('message', (topic, message) => {
console.log(`${topic}: ${message.toString()}`);
});
client.on('error', (err) => {
console.error('Connection error:', err);
});
client.on('close', () => {
console.log('Connection closed');
});
// Publish
client.publish('sensors/temp', '22.5', { qos: 1 }, (err) => {
if (err) console.error('Publish error:', err);
});
Go (paho.mqtt.golang):
package main
import (
"encoding/json"
"fmt"
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
func main() {
// Message handler
messageHandler := func(client mqtt.Client, msg mqtt.Message) {
fmt.Printf("%s: %s\n", msg.Topic(), msg.Payload())
}
// Connection handler
connectHandler := func(client mqtt.Client) {
fmt.Println("Connected")
// Subscribe on connect
if token := client.Subscribe("sensors/#", 1, nil); token.Wait() && token.Error() != nil {
fmt.Println(token.Error())
}
}
// Lost connection handler
lostHandler := func(client mqtt.Client, err error) {
fmt.Printf("Connection lost: %v\n", err)
}
// Client options
opts := mqtt.NewClientOptions()
opts.AddBroker("tcp://localhost:1883")
opts.SetClientID("go-client")
opts.SetUsername("user")
opts.SetPassword("pass")
opts.SetKeepAlive(60 * time.Second)
opts.SetDefaultPublishHandler(messageHandler)
opts.OnConnect = connectHandler
opts.OnConnectionLost = lostHandler
opts.SetAutoReconnect(true)
opts.SetMaxReconnectInterval(10 * time.Second)
// Last Will
willPayload, _ := json.Marshal(map[string]bool{"online": false})
opts.SetWill("clients/go-client/status", string(willPayload), 1, true)
// Connect
client := mqtt.NewClient(opts)
if token := client.Connect(); token.Wait() && token.Error() != nil {
panic(token.Error())
}
// Publish
payload := map[string]float64{"temperature": 22.5}
data, _ := json.Marshal(payload)
token := client.Publish("sensors/temp", 1, false, data)
token.Wait()
// Keep running
select {}
}
Examples
Command-Line Tools:
# Mosquitto clients
# Subscribe with verbose output
mosquitto_sub -h broker.example.com -p 1883 \
-u user -P pass \
-t "sensors/#" -v -q 1
# Publish with file contents
mosquitto_pub -h localhost -t "config/device" \
-f /path/to/config.json -r -q 1
# Debug mode
mosquitto_sub -t "#" -v -d
# MQTT.js CLI
npx mqtt sub -t 'sensors/#' -h 'localhost' -v
npx mqtt pub -t 'test' -m 'hello' -h 'localhost'
# MQTTX CLI (cross-platform)
mqttx sub -t 'sensors/#' -h localhost
mqttx pub -t 'test' -m 'hello' -h localhost
mqttx bench pub -t 'test' -c 100 -m 'benchmark'
GUI Tools:
- MQTT Explorer: Cross-platform, topic tree visualisation
- MQTTX: Modern UI, scriptable, benchmarking
- HiveMQ Websocket Client: Browser-based testing
Common Use Cases
Key Concepts
graph TB
subgraph "IoT Sensor Network"
S1[Temperature<br/>Sensor] -->|Publish| B1[Broker]
S2[Humidity<br/>Sensor] -->|Publish| B1
S3[Motion<br/>Sensor] -->|Publish| B1
B1 -->|Subscribe| DB[(Time Series<br/>Database)]
B1 -->|Subscribe| Alert[Alert<br/>Service]
B1 -->|Subscribe| Dash[Dashboard]
end
subgraph "Home Automation"
App[Mobile App] -->|Command| B2[Broker]
B2 -->|Control| Light[Smart Light]
B2 -->|Control| Thermo[Thermostat]
B2 -->|Status| App
end
subgraph "Fleet Management"
V1[Vehicle 1] -->|GPS/Telemetry| B3[Broker]
V2[Vehicle 2] -->|GPS/Telemetry| B3
B3 -->|Track| Monitor[Monitoring<br/>System]
B3 -->|Dispatch| Dispatch[Dispatch<br/>System]
end
Examples
IoT Sensor Data Collection:
import paho.mqtt.client as mqtt
import json
import time
import random
class IoTSensor:
def __init__(self, sensor_id, broker_host):
self.sensor_id = sensor_id
self.client = mqtt.Client(client_id=f"sensor-{sensor_id}")
# Configure LWT
self.client.will_set(
f"sensors/{sensor_id}/status",
"offline",
qos=1,
retain=True
)
self.client.connect(broker_host, 1883, keepalive=30)
self.client.loop_start()
# Announce online
self.client.publish(
f"sensors/{sensor_id}/status",
"online",
qos=1,
retain=True
)
def publish_reading(self, temperature, humidity):
payload = {
"sensor_id": self.sensor_id,
"timestamp": time.time(),
"temperature": temperature,
"humidity": humidity
}
self.client.publish(
f"sensors/{self.sensor_id}/data",
json.dumps(payload),
qos=1
)
# Simulate sensor
sensor = IoTSensor("temp-01", "localhost")
while True:
temp = 20 + random.uniform(-2, 2)
humidity = 50 + random.uniform(-10, 10)
sensor.publish_reading(temp, humidity)
time.sleep(5)
Home Automation Control:
import paho.mqtt.client as mqtt
import json
class SmartDevice:
def __init__(self, device_id, broker_host):
self.device_id = device_id
self.client = mqtt.Client(client_id=device_id)
self.client.on_connect = self._on_connect
self.client.on_message = self._on_message
self.state = {"power": False, "brightness": 100}
# LWT
self.client.will_set(
f"devices/{device_id}/availability",
"offline",
qos=1,
retain=True
)
self.client.connect(broker_host, 1883)
def _on_connect(self, client, userdata, flags, rc):
# Subscribe to commands
client.subscribe(f"devices/{self.device_id}/set/#", qos=1)
# Publish availability
client.publish(
f"devices/{self.device_id}/availability",
"online",
qos=1,
retain=True
)
# Publish initial state
self._publish_state()
def _on_message(self, client, userdata, msg):
command = msg.topic.split('/')[-1]
value = msg.payload.decode()
if command == "power":
self.state["power"] = value.lower() == "on"
elif command == "brightness":
self.state["brightness"] = int(value)
self._publish_state()
def _publish_state(self):
self.client.publish(
f"devices/{self.device_id}/state",
json.dumps(self.state),
qos=1,
retain=True
)
def run(self):
self.client.loop_forever()
# Smart light device
light = SmartDevice("living-room-light", "localhost")
light.run()
Chat/Messaging Application:
import paho.mqtt.client as mqtt
import json
import time
class ChatClient:
def __init__(self, username, broker_host):
self.username = username
self.client = mqtt.Client(client_id=f"chat-{username}")
self.client.on_connect = self._on_connect
self.client.on_message = self._on_message
# LWT for presence
self.client.will_set(
f"chat/presence/{username}",
json.dumps({"online": False}),
qos=1,
retain=True
)
self.client.connect(broker_host, 1883)
def _on_connect(self, client, userdata, flags, rc):
# Subscribe to messages
client.subscribe([
(f"chat/room/+", 1), # Room messages
(f"chat/direct/{self.username}", 1), # Direct messages
("chat/presence/+", 1) # Presence updates
])
# Announce online
client.publish(
f"chat/presence/{self.username}",
json.dumps({"online": True, "timestamp": time.time()}),
qos=1,
retain=True
)
def _on_message(self, client, userdata, msg):
data = json.loads(msg.payload)
if "chat/room/" in msg.topic:
room = msg.topic.split('/')[-1]
print(f"[{room}] {data['from']}: {data['message']}")
elif "chat/direct/" in msg.topic:
print(f"[DM] {data['from']}: {data['message']}")
elif "chat/presence/" in msg.topic:
user = msg.topic.split('/')[-1]
status = "online" if data.get("online") else "offline"
print(f"* {user} is {status}")
def send_room(self, room, message):
self.client.publish(
f"chat/room/{room}",
json.dumps({"from": self.username, "message": message}),
qos=1
)
def send_direct(self, to_user, message):
self.client.publish(
f"chat/direct/{to_user}",
json.dumps({"from": self.username, "message": message}),
qos=1
)
def run(self):
self.client.loop_forever()
Fleet Tracking:
import paho.mqtt.client as mqtt
import json
class FleetMonitor:
def __init__(self, broker_host):
self.client = mqtt.Client(client_id="fleet-monitor")
self.client.on_message = self._on_message
self.vehicles = {}
self.client.connect(broker_host, 1883)
self.client.subscribe([
("fleet/+/gps", 1),
("fleet/+/status", 1),
("fleet/+/diagnostics", 1)
])
def _on_message(self, client, userdata, msg):
parts = msg.topic.split('/')
vehicle_id = parts[1]
data_type = parts[2]
data = json.loads(msg.payload)
if vehicle_id not in self.vehicles:
self.vehicles[vehicle_id] = {}
self.vehicles[vehicle_id][data_type] = data
# Check for alerts
if data_type == "gps":
self._check_geofence(vehicle_id, data)
elif data_type == "diagnostics":
self._check_diagnostics(vehicle_id, data)
def _check_geofence(self, vehicle_id, gps_data):
# Check if vehicle is outside allowed area
pass
def _check_diagnostics(self, vehicle_id, diag_data):
if diag_data.get("fuel_level", 100) < 20:
self.client.publish(
f"alerts/fleet/{vehicle_id}",
json.dumps({"type": "low_fuel", "level": diag_data["fuel_level"]}),
qos=1
)
def run(self):
self.client.loop_forever()
Quick Reference
| Category | Command/Pattern | Description |
|---|---|---|
| Subscribe | mosquitto_sub -t "topic/#" -v |
Subscribe with verbose output |
| Publish | mosquitto_pub -t "topic" -m "msg" |
Publish message |
| Retained | mosquitto_pub -t "topic" -m "msg" -r |
Publish retained message |
| QoS 1 | mosquitto_pub -t "topic" -m "msg" -q 1 |
Publish with at-least-once delivery |
| QoS 2 | mosquitto_pub -t "topic" -m "msg" -q 2 |
Publish with exactly-once delivery |
| Clear Retained | mosquitto_pub -t "topic" -m "" -r -n |
Delete retained message |
| LWT | --will-topic "t" --will-payload "p" |
Set Last Will message |
| Wildcard + | sensors/+/temp |
Match single level |
| Wildcard # | sensors/# |
Match multiple levels |
| TLS | --cafile ca.crt --cert client.crt |
Connect with TLS |
| Auth | -u username -P password |
Authenticate to broker |
Common Issues and Solutions
Connection Issues
| Issue | Cause | Solution |
|---|---|---|
Connection refused |
Broker not running or wrong port | Verify broker is running: systemctl status mosquitto |
Not authorised |
Invalid credentials or ACL | Check username/password and ACL configuration |
Connection lost |
Keepalive timeout | Increase keepalive interval or improve network |
Client ID in use |
Duplicate client ID | Use unique client IDs or set clean_session=True |
Message Issues
| Issue | Cause | Solution |
|---|---|---|
| Messages not received | Wrong topic or QoS mismatch | Verify topic subscription and wildcards |
| Duplicate messages | QoS 1 retry behaviour | Implement idempotent message handling |
| Retained message not received | No retained message set | Publish with retain flag first |
| Messages delayed | QoS 2 handshake | Use QoS 0/1 for time-sensitive data |
Performance Issues
| Issue | Cause | Solution |
|---|---|---|
| High latency | QoS 2 overhead or network | Use QoS 0/1, optimise network |
| Broker overload | Too many connections | Implement connection pooling, use clustering |
| Memory exhaustion | Large message queues | Set max_queued_messages, increase persistence |
| Slow subscribers | Subscriber can't keep up | Implement backpressure, filter messages |
Debugging Commands
# Check broker status
systemctl status mosquitto
mosquitto -v # Verbose mode
# Test connectivity
mosquitto_pub -h broker -t "test" -m "ping" -d
mosquitto_sub -h broker -t "test" -d
# Monitor all traffic (debug)
mosquitto_sub -t "#" -v
# Check broker statistics
mosquitto_sub -t "\$SYS/#" -v
# View active connections (if enabled)
mosquitto_sub -t "\$SYS/broker/clients/connected"
# Verify TLS
openssl s_client -connect broker:8883 -CAfile ca.crt
Python Debugging
import paho.mqtt.client as mqtt
import logging
# Enable debug logging
logging.basicConfig(level=logging.DEBUG)
logger = logging.getLogger(__name__)
client = mqtt.Client()
client.enable_logger(logger)
# Connection debugging
def on_connect(client, userdata, flags, rc):
if rc != 0:
reasons = {
1: "Incorrect protocol version",
2: "Invalid client identifier",
3: "Server unavailable",
4: "Bad username or password",
5: "Not authorised"
}
logger.error(f"Connection failed: {reasons.get(rc, 'Unknown')}")
def on_log(client, userdata, level, buf):
logger.debug(f"MQTT: {buf}")
client.on_connect = on_connect
client.on_log = on_log
Related Topics
- Message Queue Patterns: Design patterns for reliable messaging, dead letter queues, and message routing
- WebSockets: Alternative real-time communication for browser-based applications
- Apache Kafka: High-throughput distributed streaming for large-scale data pipelines
- InfluxDB/TimescaleDB: Time-series databases for storing MQTT sensor data
- Node-RED: Visual flow-based programming for IoT applications with MQTT integration
- Home Assistant: Home automation platform with extensive MQTT device support