KubeMQ
ConnectorsMQTTScenarios

Migrating from MQTT

Move an MQTT 3.1.1 or 5.0 app to KubeMQ by changing the broker host: quick start, topic-prefix mapping, persistence, and what does not migrate.

You migrate an existing MQTT application to KubeMQ by changing only the broker endpoint — the host:port your client already connects to. This is a drop-in, endpoint-only migration: no KubeMQ SDK and no library swap. Your standard MQTT 3.1.1 / 5.0 client keeps publishing and subscribing exactly as it does today. The one thing to plan for is KubeMQ's topic grammar — the connector maps MQTT topics onto four messaging patterns through a topic-prefix convention, and a handful of common MQTT features are intentionally unsupported. Read this guide before cutting over so you know what works, what needs a topic rename, and what does not migrate at all.

Overview

KubeMQ ships an embedded MQTT broker that speaks MQTT 3.1.1 and MQTT 5.0 over TCP, TLS, and WebSocket. Existing MQTT clients point at KubeMQ by changing the broker host — no library swap required. The connector maps MQTT topics onto KubeMQ's messaging patterns through a topic-prefix convention, and several common MQTT features are intentionally unsupported.

  • Default ports — 1883 (TCP), 8883 (TLS), 8083 (WebSocket — serves ws or wss depending on the Security config).
  • Canonical client — Eclipse Paho paho-mqtt 2.x (Python). The examples in this guide use it; any conformant MQTT 3.1.1 / 5.0 client works.
  • Protocol versions — MQTT 3.1.1 and MQTT 5.0. MQTT 3.1 (protocol level 3) is rejected at CONNECT.
  • Enable — CONNECTORSMQTT_ENABLE=true (no underscore after CONNECTORS). The misspelling CONNECTORS_MQTT_ENABLE is ignored, and the startup log warns about it. On Kubernetes: spec.mqtt.enabled: true.

Quick start: repoint one client

Try the switch on your own machine before you plan the migration.

The Docker tab points one client at the KubeMQ you started with Try KubeMQ, turning its connector on where needed and keeping your evaluation or license and your data. The Kubernetes tab works on a cluster you installed with Install on Kubernetes. Use it to try the switch, not for production; see Plans compared for staging and production licenses. If your current broker runs on this machine, stop it first, because two programs cannot publish the same port. The kmq command runs natively on Windows; the other commands are bash, so on Windows run them in WSL (Windows Subsystem for Linux).

You started KubeMQ with Try KubeMQ or Install with Docker.

Turn the MQTT connector on

The connector is off by default. This command turns it on and publishes port 1883 on the same container and volume. If you started KubeMQ with the Docker command, skip it and use Without kmq below.

Terminal
kmq deploy update --installation kubemq --connector mqtt --deadline 10m

Check that installation_id is the same as before and data_retained is true. You should see:

Output
{
  "installation_id": "…",
  "state": "…",
  "target": "…",
  "data_retained": true,
  "next_action": "…"
}

Point your client at KubeMQ

Before
mqtt://broker.example.com:1883
After
mqtt://localhost:1883

Topics without a reserved prefix go to Events, which a disconnected subscriber never receives. See Prefixless topics are not stored.

Send and receive one message

Install the Mosquitto clients (brew install mosquitto on macOS, sudo apt-get install mosquitto-clients on Debian or Ubuntu).

Start a subscriber in a second terminal and leave it running:

Terminal 2
mosquitto_sub -h localhost -p 1883 -t 'events/smoke/test'

Publish one message:

Terminal
mosquitto_pub -h localhost -p 1883 -t 'events/smoke/test' -m '{"ok":true}'

You should see, in Terminal 2:

Output
{"ok":true}

To turn the connector off again, run the same kmq deploy update command without --connector. Your messages and your evaluation or license stay.

You installed KubeMQ with Install on Kubernetes.

Turn the connector on

Put the connector's key in cluster-values.yaml. If Helm installed your cluster, add it to the file you created on Install on Kubernetes; if kmq installed it, create the file with only this key.

cluster-values.yaml (add)
mqtt:
  enabled: true

If kmq installed your cluster, apply the file with kmq:

Terminal
kmq deploy update --installation messaging --kube-context YOUR_KUBE_CONTEXT --values cluster-values.yaml --deadline 10m

Replace:

  • YOUR_KUBE_CONTEXT — your cluster's kubectl context; kubectl config get-contexts lists them.

Check that data_retained is true. You should see:

Output
{
  "installation_id": "…",
  "state": "…",
  "target": "…",
  "data_retained": true,
  "next_action": "…"
}

If Helm installed your cluster, apply the file with the command in Apply a settings change.

Either way, check that the connector's Service exists:

Terminal
kubectl --context YOUR_KUBE_CONTEXT -n kubemq get svc messaging-mqtt

Check that the row lists 1883/TCP. You should see:

Output
NAME     TYPE        CLUSTER-IP   EXTERNAL-IP   PORT(S)          AGE
messaging-mqtt    ClusterIP   …            <none>        1883/TCP,…     …

Reach the connector from your machine

Keep this running in a second terminal:

Terminal 2
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-mqtt 1883:1883

You should see:

Output
Forwarding from 127.0.0.1:1883 -> 1883
Forwarding from [::1]:1883 -> 1883

Send and receive one message

Run steps 2 and 3 of the Docker tab, using a third terminal for the command the Docker tab runs in Terminal 2. The port-forward serves the connector at the same localhost:1883 address.

Before you move production traffic

  • Read What does not migrate before you commit.
  • Plan how traffic moves and how you roll back: run both brokers, move one topic tree at a time, and keep your MQTT broker until you have verified the move.
  • Install KubeMQ for real: Choose your path.

Compatibility Matrix

This matrix summarizes what the MQTT connector supports. It is self-contained — read the footnotes for the load-bearing caveats.

DimensionSupportNotes
Drop-in levelendpoint-only⁵Host-swap only; topic names must follow the prefix convention
Point-to-point queues✅queues/* topics → KubeMQ Queues; $share/{g}/queues/* for consumption
Pub/sub (non-durable)✅events/* and prefixless topics (when DefaultPattern=events)
Durable / persistent subscriptions⚠️store/* → Events Store; subscriptions are always StartNewOnly — no replay
Request/reply (RPC)✅ (v5 only)MQTT 5.0 Response Topic + Correlation Data → Commands/Queries; MQTT 3.1.1 lacks Response Topic — RPC not available
Ordering guarantee⚠️ QoS-dependentQoS 2 de-duplicates at the wire level; end-to-end ordering is node-local
TransactionsN/AMQTT has no transaction concept
Dead-letter / redrive❌ no client DLQ⁶Queue messages past MaxReceiveCount are silently dropped; no consumable dead-letter address over MQTT
Selectors / filtering / wildcards⚠️ wildcards Events-only+ / # wildcards supported on events/* subscriptions only; rejected (SUBACK 0xA2) on store/, queues/, commands/, queries/
Auth modelJWT (password)CONNECT password = KubeMQ JWT when auth is enabled; username is audit-only
TLS / mTLS✅ 8883 / 8083Server TLS and mutual TLS on 8883; WebSocket 8083 serves wss when Security is configured (same port as plain ws)
Top unsupportedMQTT 3.1; retained messages; clients-as-responders; wildcards on non-Events; node-local sessions

Footnote ⁵: MQTT 3.1 clients (protocol level 3) are rejected at CONNECT. Use 3.1.1 (Paho default, protocol level 4) or 5.0.

Footnote ⁶: There is no client-settable DLQ over this protocol. See No client-settable DLQ on queue subscriptions below for details.

Connection / Endpoint Migration

Replace the broker host and port. The MQTT protocol and your client library stay the same. No code changes are needed for the connection itself — only topic names may need renaming (see Concept & Destination Mapping below).

Broker endpoints
# Before (any MQTT broker)
mqtt://broker.example.com:1883
mqtts://broker.example.com:8883

# After (KubeMQ)
mqtt://kubemq.example.com:1883
mqtts://kubemq.example.com:8883

# WebSocket — single listener on 8083
# Serves plain ws:// when Security is not configured,
# or wss:// (TLS) on the same port when Security is configured.
ws://kubemq.example.com:8083/    # plain (no Security block)
# wss://kubemq.example.com:8083/ # TLS (when Security is configured)

Protocol version requirement

The connector enforces MinProtocolVersion = 4 (MQTT 3.1.1) by default. MQTT 3.1 clients (protocol level 3) are rejected at CONNECT. Eclipse Paho 2.x defaults to MQTT 3.1.1, so no version flag is needed unless you want MQTT 5.0.

# paho-mqtt 2.x — explicitly request MQTT 5.0 (optional; 3.1.1 is the default)
import paho.mqtt.client as mqtt

client = mqtt.Client(
    callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
    protocol=mqtt.MQTTv5,   # omit for 3.1.1 (the Paho default)
)

Authentication

When the KubeMQ authentication service is enabled, the CONNECT password must be a valid KubeMQ JWT. The username field is accepted and recorded for audit but is not validated.

client.username_pw_set(username="device-01", password="<kubemq-jwt>")

When authentication is disabled (the default), all connections are accepted regardless of credentials — exactly like the gRPC and REST listeners. If the MQTT ports are reachable from untrusted networks, either enable authentication or restrict access at the network layer.

Concept & Destination Mapping

The first segment of every MQTT topic selects a KubeMQ messaging pattern. The remainder of the topic path (after the prefix slash) becomes the KubeMQ channel name, with / translated to ..

MQTT topicKubeMQ patternKubeMQ channelNotes
events/sensor/tempEventssensor.tempPub/sub, at-most-once
store/audit/loginEvents Storeaudit.loginPersistent, StartNewOnly — no replay
queues/orders/newQueuesorders.newAt-least-once; consume via $share
commands/device/rebootCommandsdevice.rebootRPC; MQTT 5.0 only
queries/inventory/statusQueriesinventory.statusRPC; MQTT 5.0 only
sensor/tempDefaultPatternsensor.tempPrefixless topic routed by config

Prefix collision caution. If your existing topics already use a top-level segment named events, store, queues, commands, or queries, those topics will be interpreted as KubeMQ pattern prefixes rather than literal topic names. Rename those segments before migrating.

Prefixless topics

Topics whose first segment is not a reserved prefix are routed according to DefaultPattern (default: events), so basic pub/sub keeps working with no topic renames. Set DefaultPattern = none to reject prefixless topics explicitly.

Wildcard subscriptions

KubeMQ supports the two standard MQTT wildcards on events/ subscriptions only:

  • + — single-level wildcard (matches exactly one topic segment).
  • # — multi-level wildcard (matches the rest of the topic tree).

Wildcard use on any other pattern is rejected; exact-filter subscriptions work on all patterns. See Wildcards on non-Events patterns — REJECTED for the exact SUBACK codes and full behavior.

Queue consumption (shared subscriptions)

Plain subscriptions to queues/* are rejected (SUBACK 0x83). Queue consumption requires a shared subscription with QoS ≥ 1:

$share/{group}/queues/{path}

Example: $share/workers/queues/orders/new

The group name is recorded for metrics and audit but does not create independent per-group message copies — all groups and all nodes compete in the same KubeMQ queue pool.

/ and . conflation in channel names

KubeMQ uses . as the channel hierarchy separator. The topic mapper replaces / with ., so both devices/a/b/temp and devices/a.b.temp map to the same channel devices.a.b.temp. Design channel names to avoid ambiguity if you have topics that mix both separators.

Canonical Client Example

Client: Eclipse Paho paho-mqtt 2.x

The snippets below cover the core flows you migrate: publish an event, subscribe and consume events, consume a queue via $share, an MQTT 5.0 RPC request, and a TLS connection.

Install

Terminal
pip install paho-mqtt>=2.0.0

Publish an event

import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, reason_code, properties):
    print(f"Connected: reason_code={reason_code}")

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
# When KubeMQ authentication is enabled, set password = KubeMQ JWT:
# client.username_pw_set(username="device-01", password="<kubemq-jwt>")
client.on_connect = on_connect

client.connect("kubemq.example.com", 1883)
client.loop_start()

# Publish to Events pattern → channel "sensor.temp"
result = client.publish("events/sensor/temp", payload='{"celsius": 21.5}', qos=1)
result.wait_for_publish()

client.loop_stop()
client.disconnect()

Subscribe and consume events

import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, reason_code, properties):
    print(f"Connected: reason_code={reason_code}")
    # Subscribe to Events with a single-level '+' wildcard (Events-only)
    client.subscribe("events/sensor/+", qos=1)

def on_message(client, userdata, message):
    print(f"Topic: {message.topic}  Payload: {message.payload.decode()}")

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
client.on_connect = on_connect
client.on_message = on_message

client.connect("kubemq.example.com", 1883)
client.loop_forever()

Consume a queue (at-least-once)

Queue consumption requires a shared subscription and QoS ≥ 1. The PUBACK doubles as the queue acknowledge — the message is requeued if the PUBACK is not sent within QueueAckTimeoutSeconds (default 30 s).

import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, reason_code, properties):
    print(f"Connected: reason_code={reason_code}")
    # $share group name is audit-only; all groups share the same queue pool
    client.subscribe("$share/workers/queues/orders/new", qos=1)

def on_message(client, userdata, message):
    print(f"Queue message: {message.payload.decode()}")
    # paho-mqtt 2.x auto-sends PUBACK for QoS 1 when on_message returns.
    # The KubeMQ connector treats PUBACK receipt as the queue acknowledgment.

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
client.on_connect = on_connect
client.on_message = on_message

client.connect("kubemq.example.com", 1883)
client.loop_forever()

RPC (MQTT 5.0 only)

RPC requires MQTT 5.0 (Response Topic + Correlation Data). MQTT 3.1.1 does not support RPC on this connector — see RPC over MQTT 3.1.1 — NOT AVAILABLE for the exact behavior.

import paho.mqtt.client as mqtt
from paho.mqtt.properties import Properties
from paho.mqtt.packettypes import PacketTypes
import uuid

CLIENT_ID = "req-client-01"
REPLY_TOPIC = f"$reply/{CLIENT_ID}/inventory"

def on_connect(client, userdata, flags, reason_code, properties):
    print(f"Connected: reason_code={reason_code}")
    # Subscribe to our own $reply inbox before sending the request
    client.subscribe(REPLY_TOPIC, qos=1)

def on_message(client, userdata, message):
    print(f"RPC response: {message.payload.decode()}")
    # Check user properties for command outcome:
    #   kubemq-executed = "true"/"false"  (Commands)
    #   kubemq-metadata, kubemq-error    (Queries / error detail)

client = mqtt.Client(
    callback_api_version=mqtt.CallbackAPIVersion.VERSION2,
    client_id=CLIENT_ID,
    protocol=mqtt.MQTTv5,
)
client.on_connect = on_connect
client.on_message = on_message
client.connect("kubemq.example.com", 1883)
client.loop_start()

# Build MQTT 5.0 publish properties
pub_props = Properties(PacketTypes.PUBLISH)
pub_props.ResponseTopic = REPLY_TOPIC
pub_props.CorrelationData = str(uuid.uuid4()).encode()

# The responder is a gRPC/REST/CloudEvents client subscribed to channel "inventory.status"
client.publish(
    "queries/inventory/status",
    payload='{"sku":"WIDGET-100"}',
    qos=1,
    properties=pub_props,
)

import time; time.sleep(5)   # wait for response
client.loop_stop()
client.disconnect()

TLS connection

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
client.tls_set(
    ca_certs="/path/to/ca.crt",       # CA that signed the KubeMQ server cert
    certfile="/path/to/client.crt",   # omit for server-TLS-only (no mTLS)
    keyfile="/path/to/client.key",
)
client.connect("kubemq.example.com", 8883)

For WebSocket TLS: client.connect("kubemq.example.com", 8083, transport="websockets") with the same tls_set call above. Port 8083 serves wss:// when Security is configured (the same port that serves ws:// without TLS).

Security

Authentication

Authentication is connect-time only. A JWT that expires mid-connection does not terminate the connection — the token is validated once at CONNECT and cached for the connection's lifetime.

When auth is ...Behavior
Disabled (default)All CONNECT attempts accepted; username recorded for audit
EnabledCONNECT password must be a valid KubeMQ JWT; empty or invalid → CONNACK 0x86 (MQTT 5.0) / rc=5 (3.1.1)

Authorization

When the KubeMQ authorization service (Casbin) is configured, every publish and subscribe is checked against the client's resolved identity. ACL denial returns PUBACK 0x87 / SUBACK 0x87 (MQTT 5.0) or silent drop / SUBACK 0x80 (MQTT 3.1.1); the $reply/<own-client-id>/... namespace is always allowed.

TLS / mTLS

The dedicated TLS listener on 8883 activates only when the global Security block is configured with certificate material; without it, 8883 is skipped at startup with a warning. The WebSocket listener on 8083 is a single listener that always starts when WsPort is set: it serves plain ws:// when Security is not configured, and switches to wss:// (TLS) on the same port when Security is configured. The plain TCP listener on 1883 starts unconditionally when the connector is enabled.

  • Server TLS — Security.Mode = tls: presents Security.Cert / Security.Key; clients verify the server certificate; minimum TLS 1.2.
  • Mutual TLS — Security.Mode = mtls: additionally requires clients to present a certificate verified against Security.Ca.

Opt-in activation

The connector is disabled by default. Explicitly enable it:

config.toml
[Connectors.MQTT]
  Enable = true
  Port   = "1883"

Or via environment variable:

Terminal
export CONNECTORSMQTT_ENABLE=true

What Does NOT Migrate / Documented Deviations

Prefixless topics are not stored

DefaultPattern is events, so a topic without a reserved prefix is an Events channel. Events are not stored: a subscriber that is disconnected when a message is published never receives it, at any QoS. For messages that must wait for a consumer, publish to a queues/ topic and consume it with a $share subscription (see Queue consumption). A store/ topic persists messages in the Events Store, but an MQTT subscriber receives only messages published after it subscribes. CONNECTORSMQTT_DEFAULT_PATTERN=store routes every prefixless topic to the Events Store.

MQTT 3.1 (protocol level 3) — REJECTED

Clients using MQTT 3.1 (protocol level 3, e.g. very old Paho 1.x defaults) receive CONNACK with "unsupported protocol version". Upgrade to 3.1.1 (mqtt.MQTTv311) or 5.0 (mqtt.MQTTv5). Paho 2.x defaults to 3.1.1 — no change needed for modern clients.

Retained messages — REJECTED AND AUDITED (not silently dropped)

RetainAvailable is forced to 0 at server construction. Any publish with the retain flag set is rejected with an audit event (publish.error) and is not routed and not stored. The MQTT 5.0 client receives a normal success PUBACK (the retain flag on a PUBACK has no standard meaning in MQTT), but the message is discarded before reaching KubeMQ. MQTT 3.1.1 clients also have their retained publishes discarded with an audit event. There is no path to convert retained messages to Events Store entries.

A Will (LWT) message with the retain flag set is rejected at CONNECT (CONNACK 0x9A), since retain is unavailable. Use an LWT without the retain flag.

MQTT clients as RPC responders — NOT SUPPORTED

A SUBSCRIBE to commands/* or queries/* is rejected with SUBACK 0x83. MQTT clients can only issue RPC requests (as publishers); RPC responders must be gRPC, REST, or CloudEvents clients. This is an architectural constraint.

RPC over MQTT 3.1.1 — NOT AVAILABLE

MQTT 3.1.1 has no Response Topic field. A publish to commands/* or queries/* from a 3.1.1 client is silently dropped and audited (publish.error + WARN). RPC is exclusively an MQTT 5.0 capability on this connector.

Wildcards on non-Events patterns — REJECTED

Wildcard subscriptions (+, #) are supported only on the events/ prefix (and on prefixless topics when DefaultPattern=events). Attempts to use wildcards on store/, queues/, commands/, or queries/ receive SUBACK 0xA2 (MQTT 5.0) or SUBACK 0x80 (MQTT 3.1.1). Exact-filter subscriptions work on all patterns.

Sessions — node-local, in-memory, lost on restart

MQTT session state (subscriptions, client-ID registry, LWT, QoS state machines) is held in memory on the node that accepted the connection. Consequences:

  • Process restart clears all session state, including clean-session=0 sessions.
  • A clean-session=0 client that reconnects to a different cluster node finds no session.
  • Duplicate client-ID from different nodes keeps both connections live (no cross-node takeover).
  • Cross-node RPC responses are lost: if the requester's connection migrates to another node mid-request, the response is not forwarded.

Events Store replay — not available over MQTT

Subscriptions to store/* are always StartNewOnly. Historical replay from the Events Store is not exposed over the MQTT wire protocol. Replay is available via gRPC, REST, and CloudEvents clients only.

No client-settable DLQ on queue subscriptions

The connector never sets a per-message receive limit on messages it publishes to KubeMQ Queues. A "poison" message that exceeds the server's MaxReceiveCount is dropped by the server; there is no consumable dead-letter address accessible over MQTT. See the AWS or RabbitMQ guides for connectors that support DLQ/redrive.

$SYS topics

The $SYS broker-internal topic tree is not available. Subscriptions to $SYS/* are denied.

Events Store subscriptions and the $share prefix

$share subscriptions are valid only for the queues/* pattern. $share on store/*, events/*, commands/*, or queries/* is rejected with SUBACK 0x83.

User properties (MQTT 5.0 ↔ KubeMQ Tags)

MQTT 5.0 user properties map to KubeMQ message tags and vice versa. Caps: 32 properties and 4096 bytes total. Exceeding either cap returns PUBACK 0x97. MQTT 3.1.1 clients do not support user properties — tag metadata is not surfaced to 3.1.1 subscribers.

Verification Smoke Test

For the basic round trip, see the quick start; this section adds deeper checks with paho-mqtt 2.x, the canonical client for this guide.

Step 1 — Publish and receive an event

Run these two snippets in separate terminals.

Terminal 1 — subscriber:

import paho.mqtt.client as mqtt

received = []

def on_connect(client, userdata, flags, rc, props):
    client.subscribe("events/smoke/test", qos=1)

def on_message(client, userdata, msg):
    received.append(msg.payload.decode())
    print(f"RECEIVED: {msg.payload.decode()}")
    client.disconnect()

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
client.on_connect = on_connect
client.on_message = on_message
client.connect("localhost", 1883)
client.loop_forever()
assert received == ['{"ok":true}'], f"Expected message not received: {received}"
print("Smoke test PASSED")

Terminal 2 — publisher (after the subscriber is connected):

import paho.mqtt.client as mqtt

client = mqtt.Client(callback_api_version=mqtt.CallbackAPIVersion.VERSION2)
client.connect("localhost", 1883)
client.loop_start()
result = client.publish("events/smoke/test", payload='{"ok":true}', qos=1)
result.wait_for_publish()
client.loop_stop()
client.disconnect()
print("Published")

You should see, in the subscriber's terminal:

Output
RECEIVED: {"ok":true}
Smoke test PASSED

Step 2 — Confirm a queue round-trip

Publish to a queue:

client.publish("queues/smoke/test", payload='{"job":1}', qos=1)

Consume from the queue (shared subscription, QoS 1):

client.subscribe("$share/smokers/queues/smoke/test", qos=1)

Expected: the consumer receives {"job":1} exactly once; no redelivery unless the consumer disconnects before the PUBACK.

See Also

Was this page helpful?

On this page