Migrating from STOMP
Change the STOMP broker host to KubeMQ: quick start, every destination type mapped, STOMP 1.2 acknowledgements, no transactions or selectors.
Migrate a STOMP client application onto KubeMQ's native STOMP connector. The connector implements STOMP 1.0 / 1.1 / 1.2 over plain TCP and TLS — no broker middleware is interposed, so your STOMP client connects directly to KubeMQ. This is a drop-in, endpoint-only migration: change the host and port, keep the same STOMP client library and application code.
The guide is broker-agnostic — it applies to any STOMP client driving any STOMP broker today:
stomp.py, @stomp/stompjs, go-stomp/stomp/v3, the Ruby stomp gem, Stomp.Net, or a Spring
broker-relay front-end. Only the library API differs; the wire protocol and the migration steps are
identical.
Overview
The KubeMQ STOMP connector is a faithful STOMP 1.0 / 1.1 / 1.2 raw-TCP endpoint. The code examples on
this page target stomp.py 8.x (Python) as the canonical client; other clients connect identically.
| Item | Value |
|---|---|
| Plain TCP port | 61613 (Connectors.Stomp.Port, default "61613") |
| TLS port | 61614 (Connectors.Stomp.TlsPort, default "61614"; active only when Security ≠ none) |
| STOMP versions | 1.0, 1.1, 1.2 (negotiated; highest common wins) |
| Default pattern | events (configurable via Connectors.Stomp.DefaultPattern) |
| Canonical client | stomp.py 8.x (Python) |
| Enable env var | CONNECTORS_STOMP_ENABLE=true, with the underscore; CONNECTORSSTOMP_ENABLE is ignored without any warning. Kubernetes spec.stomp.enabled: true |
For the full wire-protocol contract (heartbeats, receipts, header/tag mapping, ack token internals, error table, metrics catalog) see Capabilities, Destination grammar, and Error frames. For configuration fields and environment variables see the configuration reference.
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 STOMP connector on
The connector is off by default. This command turns it on and publishes port 61613 on the same container and volume. If you started KubeMQ with the Docker command, skip it and use Without kmq below.
kmq deploy update --installation kubemq --connector stomp --deadline 10mCheck that installation_id is the same as before and data_retained is true. You should see:
{
"installation_id": "…",
"state": "…",
"target": "…",
"data_retained": true,
"next_action": "…"
}Point your client at KubeMQ
import stomp
conn = stomp.Connection12([("stomp-host", 61613)])
conn.connect("user", "password", wait=True)import stomp
conn = stomp.Connection12([("localhost", 61613)])
conn.connect("user", "password", wait=True) # any login and passcode while authentication is offSend and receive one message
Install the client library and save the two scripts.
python3 -m pip install stomp-pyimport stomp, time
class Listener(stomp.ConnectionListener):
def __init__(self, conn):
self._conn = conn
def on_message(self, frame):
print("RECEIVED:", frame.body)
self._conn.ack(frame.headers["ack"])
conn = stomp.Connection12([("localhost", 61613)])
conn.set_listener("", Listener(conn))
conn.connect(wait=True)
conn.subscribe(destination="/queue/smoke-test", id="s1", ack="client-individual")
time.sleep(10) # wait for a message
conn.disconnect()import stomp
conn = stomp.Connection12([("localhost", 61613)])
conn.connect(wait=True)
conn.send(destination="/queue/smoke-test", body="hello from stomp")
conn.disconnect()
print("sent")Start the consumer in a second terminal, then run the producer within 10 seconds.
python3 consumer.pypython3 producer.pyYou should see, in Terminal 2:
RECEIVED: hello from stompTo 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.
stomp:
enabled: trueIf kmq installed your cluster, apply the file with kmq:
kmq deploy update --installation messaging --kube-context YOUR_KUBE_CONTEXT --values cluster-values.yaml --deadline 10mReplace:
YOUR_KUBE_CONTEXT— your cluster's kubectl context;kubectl config get-contextslists them.
Check that data_retained is true. You should see:
{
"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:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq get svc messaging-stompCheck that the row lists 61613/TCP. You should see:
NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE
messaging-stomp ClusterIP … <none> 61613/TCP,… …Reach the connector from your machine
Keep this running in a second terminal:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-stomp 61613:61613You should see:
Forwarding from 127.0.0.1:61613 -> 61613
Forwarding from [::1]:61613 -> 61613Send 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:61613 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 destination at a time, and keep your STOMP broker until you have verified the move.
- Install KubeMQ for real: Choose your path.
Compatibility Matrix
This is the STOMP column of the cross-protocol compatibility matrix on Migrate from another broker. It is self-contained — read it as the at-a-glance answer to "what migrates?"
| Dimension | STOMP on KubeMQ |
|---|---|
| Drop-in level | endpoint-only |
| Point-to-point queues | ✅ /queue/* → Queues pattern |
| Pub/sub (non-durable) | ✅ /topic/* → Events pattern |
| Durable / persistent subscriptions | ✅ via /topic-store/* + start-from replay headers |
| Request / reply (RPC) | ✅ reply-to |
| Ordering guarantee | ⚠️ node-local |
| Transactions | ❌ rejected (BEGIN/COMMIT/ABORT → ERROR transactions not supported) |
| Dead-letter / redrive | ❌ no client-settable DLQ; see footnote ¹ |
| Selectors / filtering / wildcards | ❌ no selectors |
| Auth model | JWT (CONNECT) |
| TLS / mTLS | ✅ 61614 when Security configured |
| Top unsupported | transactions; selectors |
¹ No client-settable DLQ over STOMP. The STOMP connector never sets
MaxReceiveQueueon published messages, so a poison message that exceedsMaxReceiveCountis silently dropped by the broker — it is not delivered to any consumable dead-letter address. If you need client-facing dead-letter behaviour, use the RabbitMQ (DLX) or AWS (redrive) path instead.
Connection / Endpoint Migration
The only change required is the host and port. The STOMP protocol remains identical.
| Setting | Source broker | KubeMQ |
|---|---|---|
| Host | stomp-host (your current broker) | kubemq-host |
| Plain port | 61613 (STOMP default) | 61613 (same) |
| TLS port | 61614 | 61614 |
| Login | broker-specific username | any string (audit-only when auth disabled); or the ClientID value when auth enabled |
| Passcode | broker password | KubeMQ JWT when Authentication.Enable = true; any value when auth disabled |
stomp.py 8.x — before
import stomp
conn = stomp.Connection12([("stomp-host", 61613)])
conn.connect("user", "password", wait=True)stomp.py 8.x — after (KubeMQ)
import stomp
conn = stomp.Connection12([("kubemq-host", 61613)])
conn.connect("user", jwt_token, wait=True) # jwt_token ignored when auth disabledThe examples use stomp.Connection12: KubeMQ adds the ack header to queue deliveries only on
STOMP 1.2 connections, and stomp.py's default Connection speaks 1.1.
For TLS, use port 61614 and call conn.set_ssl(for_hosts=[("kubemq-host", 61614)]) before connect.
Concept & Destination Mapping
The STOMP destination header determines both the KubeMQ messaging pattern and the KubeMQ channel
name. Normalization strips one leading /, splits on /, joins remaining segments with ..
Prefix table
| STOMP destination prefix | Aliases | KubeMQ pattern | Channel name |
|---|---|---|---|
/queue/NAME | /queues/NAME | Queues | NAME (.-joined segments) |
/topic/NAME | /events/NAME | Events | NAME |
/topic-store/NAME | /store/NAME | Events Store | NAME |
/command/NAME | /commands/NAME | Commands (RPC) | NAME |
/query/NAME | /queries/NAME | Queries (RPC) | NAME |
/reply/ID | — | connection-local reply | not routed to any channel |
A destination whose first segment does not match any prefix maps to Connectors.Stomp.DefaultPattern
(default events). Setting DefaultPattern = none rejects those destinations with ERROR
invalid destination.
Destination examples
| STOMP destination | KubeMQ pattern | KubeMQ channel |
|---|---|---|
/queue/orders | Queues | orders |
/queue/orders/new | Queues | orders.new |
/topic/orders.created | Events | orders.created |
/topic-store/audit-log | Events Store | audit-log |
/command/payments/process | Commands | payments.process |
/query/inventory/check | Queries | inventory.check |
Wildcards
The broker's native wildcards — * (single segment) and > (tail) — are supported on SUBSCRIBE,
Events pattern only (/topic/* or /topic/orders.>). These are the STOMP destination grammar's own
wildcards; there is no MQTT-style + / #. Wildcards on SEND, on any other pattern, or a misplaced
> are rejected with ERROR invalid destination.
Durable subscriptions → Events Store + replay headers
STOMP has no built-in durable-subscription mechanism. In KubeMQ the equivalent is subscribing to a
/topic-store/ destination with a start-from replay header:
start-from value | Meaning |
|---|---|
new (default) | Deliver only messages published after this subscription |
first | Replay from the very first stored message |
last | Start from the most recently stored message |
sequence | Start at a specific sequence number (supply start-value) |
time | Start at an RFC 3339 or unix-seconds timestamp (supply start-value) |
time-delta | Start N seconds before now (supply start-value in seconds) |
Example: subscribe from the beginning of an Events Store channel:
conn.subscribe(
destination="/topic-store/audit-log",
id="sub-1",
ack="client-individual",
headers={"start-from": "first"},
)Ack modes
Set per subscription via the ack header on SUBSCRIBE.
ack value | Semantics | Applies to |
|---|---|---|
auto (default) | Broker auto-acks on enqueue success; NAcks on enqueue failure (requeues) | Queues |
client-individual | Each delivery tracked; ACK/NACK frame releases that one message | Queues |
client | Cumulative: ACK/NACK of delivery N acks/nacks all pending with order-index ≤ N | Queues |
ACK frames on Events or Events Store subscriptions are accepted as a no-op (those patterns do not track per-message delivery state).
Canonical Client Example (stomp.py 8.x)
This example uses the following stomp.py 8.x API symbols: stomp.Connection12, conn.set_listener,
conn.connect, conn.send, conn.subscribe, conn.ack, conn.disconnect,
stomp.ConnectionListener.on_message.
Publish to a queue
import stomp
KUBEMQ_HOST = "kubemq-host"
KUBEMQ_PORT = 61613
JWT_TOKEN = "..." # omit or use any string when auth is disabled
conn = stomp.Connection12([(KUBEMQ_HOST, KUBEMQ_PORT)])
conn.connect("myapp", JWT_TOKEN, wait=True)
conn.send(
destination="/queue/orders",
body="order payload",
headers={"content-type": "text/plain"},
)
conn.disconnect()Subscribe and consume from a queue (client-individual ack)
import stomp
KUBEMQ_HOST = "kubemq-host"
KUBEMQ_PORT = 61613
JWT_TOKEN = "..."
class QueueListener(stomp.ConnectionListener):
def __init__(self, conn):
self._conn = conn
def on_message(self, frame):
print("received:", frame.body)
# Acknowledge the individual message
self._conn.ack(frame.headers["ack"])
def on_error(self, frame):
print("error:", frame.headers.get("message"))
conn = stomp.Connection12([(KUBEMQ_HOST, KUBEMQ_PORT)])
conn.set_listener("", QueueListener(conn))
conn.connect("myapp", JWT_TOKEN, wait=True)
conn.subscribe(
destination="/queue/orders",
id="sub-orders",
ack="client-individual",
)
input("Press Enter to stop...\n")
conn.disconnect()Pub/sub via Events
import stomp, threading
KUBEMQ_HOST = "kubemq-host"
KUBEMQ_PORT = 61613
received = threading.Event()
class EventListener(stomp.ConnectionListener):
def on_message(self, frame):
print("event:", frame.body)
received.set()
conn = stomp.Connection12([(KUBEMQ_HOST, KUBEMQ_PORT)])
conn.set_listener("", EventListener())
conn.connect(wait=True)
conn.subscribe(destination="/topic/orders.created", id="sub-1", ack="auto")
conn.send(destination="/topic/orders.created", body="event payload")
received.wait(timeout=5)
conn.disconnect()RPC (Commands) with reply-to
The reply subscription must be created before sending the RPC request. See What Does NOT Migrate / Deviations for details.
import stomp, threading, uuid
KUBEMQ_HOST = "kubemq-host"
KUBEMQ_PORT = 61613
JWT_TOKEN = "..."
reply_event = threading.Event()
reply_body = None
class RpcListener(stomp.ConnectionListener):
def on_message(self, frame):
global reply_body
reply_body = frame.body
reply_event.set()
reply_dest = f"/reply/{uuid.uuid4()}"
conn = stomp.Connection12([(KUBEMQ_HOST, KUBEMQ_PORT)])
conn.set_listener("", RpcListener())
conn.connect("myapp", JWT_TOKEN, wait=True)
# Step 1: subscribe to the reply destination BEFORE sending the request
conn.subscribe(destination=reply_dest, id="reply-sub", ack="auto")
# Step 2: send the RPC request with reply-to and correlation-id headers
conn.send(
destination="/command/payments/process",
body="payment request",
headers={
"reply-to": reply_dest,
"correlation-id": str(uuid.uuid4()),
"timeout": "5000", # ms; capped at RpcTimeoutSeconds * 1000
},
)
reply_event.wait(timeout=10)
print("reply:", reply_body)
conn.disconnect()Security
Authentication (connect-time only). When Authentication.Enable = true, the CONNECT frame passcode
header must carry a valid KubeMQ JWT. The login header is recorded for audit purposes. Auth failure
results in ERROR authentication failed + connection close. When authentication is disabled, any
credentials are accepted — no auth is performed.
Token expiry does not terminate an established connection (connect-time-only auth, consistent with the AMQP and MQTT connectors).
Authorization. When Authorization.Enable = true, SEND enforces the Write permission and SUBSCRIBE
enforces Read on ControlRecord{Resource: <pattern>, ClientID, Channel}. Casbin policies must cover
stomp-* client IDs for the mapped channels. Reply (/reply/...) destinations are authorization-exempt
(connection-local).
No-auth exposure note. When Authentication.Enable = false (the server default), the STOMP listener
accepts any login / passcode. If the server is reachable from untrusted networks, enable authentication
or restrict access with a firewall.
TLS. Port 61614 is active only when Connectors.Stomp.TlsPort is set and Security is configured
(mode ≠ none). Plain TCP on 61613 carries no transport encryption.
Enabling the connector.
[Connectors.Stomp]
Enable = true
Port = "61613"
TlsPort = "61614"For the full auth and Casbin authorization setup, see Authentication & Security.
What Does NOT Migrate / Deviations
Hard failures (client will receive ERROR + connection close)
| Feature | Behaviour on KubeMQ |
|---|---|
STOMP transactions — BEGIN/COMMIT/ABORT frames, or any transaction header | ERROR transactions not supported + close. Remove all transaction usage before migrating. |
Message selectors — selector header on SUBSCRIBE | ERROR selectors not supported + close. Remove selector headers. KubeMQ does not support broker-side SQL92 filtering over STOMP. |
Subscribing to RPC destinations — SUBSCRIBE to /command/* or /query/* | ERROR cannot subscribe to RPC destinations. Use the reply-to flow instead. |
Invalid wildcards — wildcards on SEND, on non-Events patterns, or misplaced > | ERROR invalid destination (see the Wildcards subsection). |
Behavioural deviations
| Feature | Deviation |
|---|---|
| RPC reply-to ordering | The reply subscription (/reply/ID) must exist on the same connection before the SEND that carries reply-to. If the reply subscription is not already active, the SEND returns ERROR reply-to subscription required + close. This differs from brokers that buffer replies for a later subscription. |
| No client-settable DLQ | See footnote ¹ in the Compatibility Matrix. Poison messages exceeding MaxReceiveCount are silently dropped — there is no consumable dead-letter address over STOMP. |
| Durable subscriptions replaced by Events Store | The durable-subscription-name / activemq.subscriptionName SUBSCRIBE header is not recognized. Use /topic-store/NAME with start-from headers for persistent, replayable subscriptions. |
| Ordering is node-local | Message ordering within a channel is maintained on the receiving node; no cross-node total ordering is guaranteed in a clustered deployment. |
| STOMP-over-WebSocket | Not supported in V1. Only plain TCP (61613) and TLS (61614) listeners exist. Spring STOMP-over-WebSocket front-ends require a V2 connector upgrade. |
ActiveMQ STOMP 1.0 selector behaviour | KubeMQ rejects selectors loudly (ERROR + close) rather than ignoring them. Applications that set selector on any SUBSCRIBE must remove the header. |
/temp-queue/, /temp-topic/ destinations | Not supported. Use a /reply/{id} connection-local destination for the request/reply use case. |
Spring /app/ destination conventions | Not a broker concept — not supported. |
Verification Smoke Test
For the basic round trip, see the quick start; this section adds deeper checks.
Events Store replay check
A subscriber on /topic-store/ receives messages published before it connected. Events Store
deliveries carry no ack header, so this consumer uses ack="auto" and never calls ack.
import stomp, time
class Listener(stomp.ConnectionListener):
def on_message(self, frame):
print("RECEIVED:", frame.body)
conn = stomp.Connection12([("localhost", 61613)])
conn.set_listener("", Listener())
conn.connect(wait=True)
conn.subscribe(destination="/topic-store/smoke-test", id="s1", ack="auto",
headers={"start-from": "first"})
time.sleep(5)
conn.disconnect()Change the destination in the quick start's producer.py to /topic-store/smoke-test and run it.
When it has printed sent, run python3 replay_consumer.py. It prints RECEIVED: hello from stomp
once for every earlier publish, although it connected after the producer disconnected.
See Also
Migrate from another broker
The connector map, the cross-protocol comparison matrix, and guides for every messaging ecosystem.
Getting Started
Connect a stock STOMP client and run a publish-and-subscribe round-trip in minutes.
Architecture
The embedded STOMP server, the frame codec, version negotiation, and the destination router.
Destination Grammar
The destination grammar you are mapping your existing destinations onto.
Capabilities
The full supported / out-of-scope matrix, ack token resolution, and the wire-protocol contract.
Configuration Reference
All STOMP connector settings and environment variables, including the CONNECTORS_STOMP_ENABLE spelling.
Migrating from ActiveMQ
Route ActiveMQ onto KubeMQ by client type — JMS, STOMP, and MQTT — when STOMP is one of several access paths.
Was this page helpful?