Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

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.

MQTT ArchitectureSubscribersPublishersPublishPublishPublishSubscribeSubscribeSubscribeDashboardSensor 1MQTT BrokerSensor 2DeviceDatabaseAlert ServiceMQTT ArchitectureSubscribersPublishersPublishPublishPublishSubscribeSubscribeSubscribeDashboardSensor 1MQTT BrokerSensor 2DeviceDatabaseAlert Service

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

Example Topicshome/living-room/temperaturehome/bedroom/lightoffice/floor1/meeting-room/occupancyTopic Hierarchy/homeofficeliving-roombedroomtemperaturehumiditytemperaturelightExample Topicshome/living-room/temperaturehome/bedroom/lightoffice/floor1/meeting-room/occupancyTopic Hierarchy/homeofficeliving-roombedroomtemperaturehumiditytemperaturelight
  • 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() raises ValueError — you must pass an API version, e.g. mqtt.Client(mqtt.CallbackAPIVersion.VERSION2), and the on_connect/on_disconnect callbacks gain reason_code/properties parameters. On 1.x the calls below work as written; on 2.x add the CallbackAPIVersion argument 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

SubscriberBrokerPublisherSubscriberBrokerPublisherQoS 0 - At Most OnceQoS 1 - At Least OnceQoS 2 - Exactly OncePUBLISHPUBLISHPUBLISHPUBLISHPUBACKPUBACKPUBLISHPUBRECPUBRELPUBLISHPUBACKPUBCOMPSubscriberBrokerPublisherSubscriberBrokerPublisherQoS 0 - At Most OnceQoS 1 - At Least OnceQoS 2 - Exactly OncePUBLISHPUBLISHPUBLISHPUBLISHPUBACKPUBACKPUBLISHPUBRECPUBRELPUBLISHPUBACKPUBCOMP
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

NewSubscriberBrokerPublisherNewSubscriberBrokerPublisherStores last retained messageSubscribes laterNew messages replace retainedUpdates stored messagePUBLISH (retained=true)SUBSCRIBEPUBLISH (retained message)PUBLISH (retained=true, new value)NewSubscriberBrokerPublisherNewSubscriberBrokerPublisherStores last retained messageSubscribes laterNew messages replace retainedUpdates stored messagePUBLISH (retained=true)SUBSCRIBEPUBLISH (retained message)PUBLISH (retained=true, new value)
  • 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

SubscribersBrokerClientSubscribersBrokerClientStores LWT messageNormal operationUnexpected disconnectDetects client goneCONNECT (with LWT)CONNACKPUBLISH messagesConnection lostPUBLISH (LWT message)SubscribersBrokerClientSubscribersBrokerClientStores LWT messageNormal operationUnexpected disconnectDetects client goneCONNECT (with LWT)CONNACKPUBLISH messagesConnection lostPUBLISH (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

Fleet ManagementGPS/TelemetryGPS/TelemetryTrackDispatchVehicle 1BrokerVehicle 2MonitoringSystemDispatchSystemHome AutomationCommandControlControlStatusMobile AppBrokerSmart LightThermostatIoT Sensor NetworkPublishPublishPublishSubscribeSubscribeSubscribeTemperatureSensorBrokerHumiditySensorMotionSensorTime SeriesDatabaseAlertServiceDashboardFleet ManagementGPS/TelemetryGPS/TelemetryTrackDispatchVehicle 1BrokerVehicle 2MonitoringSystemDispatchSystemHome AutomationCommandControlControlStatusMobile AppBrokerSmart LightThermostatIoT Sensor NetworkPublishPublishPublishSubscribeSubscribeSubscribeTemperatureSensorBrokerHumiditySensorMotionSensorTime SeriesDatabaseAlertServiceDashboard

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