KubeMQ
ConnectorsSTOMPScenarios

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.

ItemValue
Plain TCP port61613 (Connectors.Stomp.Port, default "61613")
TLS port61614 (Connectors.Stomp.TlsPort, default "61614"; active only when Security ≠ none)
STOMP versions1.0, 1.1, 1.2 (negotiated; highest common wins)
Default patternevents (configurable via Connectors.Stomp.DefaultPattern)
Canonical clientstomp.py 8.x (Python)
Enable env varCONNECTORS_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.

Terminal
kmq deploy update --installation kubemq --connector stomp --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
import stomp

conn = stomp.Connection12([("stomp-host", 61613)])
conn.connect("user", "password", wait=True)
After
import stomp

conn = stomp.Connection12([("localhost", 61613)])
conn.connect("user", "password", wait=True)   # any login and passcode while authentication is off

Send and receive one message

Install the client library and save the two scripts.

Terminal
python3 -m pip install stomp-py
consumer.py
import 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()
producer.py
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.

Terminal 2
python3 consumer.py
Terminal
python3 producer.py

You should see, in Terminal 2:

Output
RECEIVED: hello from stomp

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)
stomp:
  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-stomp

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

Output
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:

Terminal 2
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-stomp 61613:61613

You should see:

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

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: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?"

DimensionSTOMP on KubeMQ
Drop-in levelendpoint-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 modelJWT (CONNECT)
TLS / mTLS✅ 61614 when Security configured
Top unsupportedtransactions; selectors

¹ No client-settable DLQ over STOMP. The STOMP connector never sets MaxReceiveQueue on published messages, so a poison message that exceeds MaxReceiveCount is 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.

SettingSource brokerKubeMQ
Hoststomp-host (your current broker)kubemq-host
Plain port61613 (STOMP default)61613 (same)
TLS port6161461614
Loginbroker-specific usernameany string (audit-only when auth disabled); or the ClientID value when auth enabled
Passcodebroker passwordKubeMQ 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 disabled

The 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 prefixAliasesKubeMQ patternChannel name
/queue/NAME/queues/NAMEQueuesNAME (.-joined segments)
/topic/NAME/events/NAMEEventsNAME
/topic-store/NAME/store/NAMEEvents StoreNAME
/command/NAME/commands/NAMECommands (RPC)NAME
/query/NAME/queries/NAMEQueries (RPC)NAME
/reply/ID—connection-local replynot 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 destinationKubeMQ patternKubeMQ channel
/queue/ordersQueuesorders
/queue/orders/newQueuesorders.new
/topic/orders.createdEventsorders.created
/topic-store/audit-logEvents Storeaudit-log
/command/payments/processCommandspayments.process
/query/inventory/checkQueriesinventory.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 valueMeaning
new (default)Deliver only messages published after this subscription
firstReplay from the very first stored message
lastStart from the most recently stored message
sequenceStart at a specific sequence number (supply start-value)
timeStart at an RFC 3339 or unix-seconds timestamp (supply start-value)
time-deltaStart 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 valueSemanticsApplies to
auto (default)Broker auto-acks on enqueue success; NAcks on enqueue failure (requeues)Queues
client-individualEach delivery tracked; ACK/NACK frame releases that one messageQueues
clientCumulative: ACK/NACK of delivery N acks/nacks all pending with order-index ≤ NQueues

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.

config.toml
[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)

FeatureBehaviour on KubeMQ
STOMP transactions — BEGIN/COMMIT/ABORT frames, or any transaction headerERROR transactions not supported + close. Remove all transaction usage before migrating.
Message selectors — selector header on SUBSCRIBEERROR 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

FeatureDeviation
RPC reply-to orderingThe 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 DLQSee 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 StoreThe 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-localMessage ordering within a channel is maintained on the receiving node; no cross-node total ordering is guaranteed in a clustered deployment.
STOMP-over-WebSocketNot 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 behaviourKubeMQ rejects selectors loudly (ERROR + close) rather than ignoring them. Applications that set selector on any SUBSCRIBE must remove the header.
/temp-queue/, /temp-topic/ destinationsNot supported. Use a /reply/{id} connection-local destination for the request/reply use case.
Spring /app/ destination conventionsNot 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.

replay_consumer.py
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

Was this page helpful?

On this page