Migrating from RabbitMQ
Point a RabbitMQ app at KubeMQ by changing the AMQP 0-9-1 connection string: quick start, users, transactions, DLX and every deviation.
RabbitMQ applications using any standard AMQP 0-9-1 client library (pika, amqp091-go,
php-amqplib, Spring AMQP, amqplib for Node.js, …) can point at KubeMQ by changing only
the connection string — no code changes — for the supported feature set. KubeMQ speaks the
RabbitMQ wire dialect, so queue.declare, exchange.declare, publisher confirms, transactions,
dead-letter exchanges, and Direct Reply-To work as-is. This is an endpoint-only drop-in:
same library, same code, same protocol, with no KubeMQ SDK to adopt.
Overview
| Connector | KubeMQ RabbitMQ connector (AMQP 0-9-1 RabbitMQ dialect) |
| Ports | 5672 (AMQP plain) / 5671 (AMQPS/TLS); 15672 for the optional management HTTP API |
| Canonical client | pika 1.x (Python) |
| Drop-in level | Endpoint-only — change the connection string; keep all client library code unchanged for the supported feature set |
| Enable default | On — Connectors.Amqp.Enable = true by default. Set CONNECTORS_AMQP_ENABLE=false (or Enable = false in TOML) to close the listener; on Kubernetes, spec.amqp.enabled: false turns it off |
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.
Publish the RabbitMQ port
The connector is on, but Try KubeMQ does not publish its port, 5672. This command publishes it 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 amqp --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
amqp://user:password@rabbitmq.example.com:5672/amqp://user:password@localhost:5672/Send and receive one message
Install the client with pip install pika, then run each file with python3.
import pika
conn = pika.BlockingConnection(pika.URLParameters("amqp://user:password@localhost:5672/"))
ch = conn.channel()
ch.queue_declare(queue="smoke-test", durable=True)
ch.basic_publish(exchange="", routing_key="smoke-test", body=b"hello-kubemq")
conn.close()
print("published ok")import pika
conn = pika.BlockingConnection(pika.URLParameters("amqp://user:password@localhost:5672/"))
ch = conn.channel()
method, props, body = ch.basic_get(queue="smoke-test", auto_ack=True)
assert body == b"hello-kubemq", f"unexpected body: {body!r}"
conn.close()
print("consume ok:", body)You should see:
published ok
consume ok: b'hello-kubemq'To undo the change, run the same kmq deploy update command without --connector. Your messages and your evaluation or license stay.
You installed KubeMQ with Install on Kubernetes. The RabbitMQ connector is on by default, so there is nothing to turn on.
Reach the connector from your machine
Keep this running in a second terminal:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-amqp 5672:5672Replace:
YOUR_KUBE_CONTEXT— your cluster's kubectl context;kubectl config get-contextslists them.
You should see:
Forwarding from 127.0.0.1:5672 -> 5672
Forwarding from [::1]:5672 -> 5672Send and receive one message
Run steps 2 and 3 of the Docker tab unchanged. The port-forward serves the connector at the same localhost:5672 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 queue or consumer at a time, and keep RabbitMQ until you have verified the move.
- Install KubeMQ for real: Choose your path.
Compatibility Matrix
The cells below are the RabbitMQ column of the master cross-protocol matrix on Migrate from another broker.
| Dimension | RabbitMQ on KubeMQ |
|---|---|
| Point-to-point queues | ✅ |
| Pub/sub (non-durable) | ✅ exchanges |
| Durable / persistent subscriptions | ✅ durable queues |
| Request / reply (RPC) | ✅ Direct Reply-To, including across cluster nodes |
| Ordering guarantee | ✅ per-queue¹ |
| Transactions | ✅ tx.select / tx.commit / tx.rollback with RabbitMQ semantics (one bound — see below) |
| Dead-letter / redrive | ✅ DLX + x-death on rejected, expired and maxlen |
| Selectors / filtering / wildcards | ✅ topic/headers routing |
| Auth model | SASL PLAIN / AMQPLAIN (KubeMQ token or connector-local users) + EXTERNAL (opt-in, mutual TLS) |
| TLS / mTLS | ✅ 5671 |
| Management HTTP API | ✅ opt-in, read-mostly subset plus definitions.json export/import — see Management HTTP API |
| Top unsupported | See Never coming, Hard rejections and Inert arguments below |
¹ Requeued messages re-enter at the queue tail, not near the head — see deviation 1 under Behavioral deviations.
Connection / Endpoint Migration
Replace the RabbitMQ host and port with the KubeMQ host. The URI scheme, client library, and application code stay the same:
# Before (RabbitMQ)
amqp://user:password@rabbitmq.example.com:5672/
amqps://user:password@rabbitmq.example.com:5671/myvhost
# After (KubeMQ)
amqp://user:password@kubemq.example.com:5672/
amqps://user:password@kubemq.example.com:5671/myvhostKey differences to be aware of at connection time:
| Aspect | RabbitMQ | KubeMQ |
|---|---|---|
| SASL mechanism | AMQPLAIN ANONYMOUS PLAIN advertised by default (measured on RabbitMQ 4.3.4); EXTERNAL once its plugin is enabled | PLAIN and AMQPLAIN always; EXTERNAL only when Connectors.Amqp.SslCertLogin = true and the listener verified a client certificate. ANONYMOUS is never offered (deviation 18) |
| Password | RabbitMQ user password | No Connectors.Amqp.Credentials configured (default): a KubeMQ JWT when Authentication.Enable = true; any value when authentication is off (username recorded for audit, never checked). Credentials configured: the RabbitMQ-style username/password pair, checked against that store |
| Vhost | Must be pre-created | Namespace-on-first-use — any charset-valid vhost is accepted; / maps to the configured DefaultVhost (default "default"). A vhost name may not contain . anywhere |
| Client identity | Connection name (optional) | Four derivations — see Which client id a policy must name. connection_name is display/log only and never reaches the identity |
| Token lifetime | n/a | A connection whose JWT expires is closed with connection.close(320); connection.update-secret refreshes it in place — see Token expiry and refresh |
| TLS | Per-listener config | Active when the server's Security block is configured; serves the same certificates as gRPC/REST |
| Queue storage | RabbitMQ queues | KubeMQ Queue channels amqp.{vhost}.{queue} |
Concept & Destination Mapping
AMQP's exchange/binding model is implemented connector-side as virtual routing. Only queues hold messages; exchanges and bindings are metadata resolved at publish time.
| RabbitMQ concept | KubeMQ pattern | Channel / address |
|---|---|---|
| Queue | Queues | amqp.{vhost}.{queue} |
Exchange (direct / fanout / topic / headers / x-delayed-message) | Virtual routing to queues | — |
| Binding (queue-to-exchange and exchange-to-exchange) | Routing rule (vhost-scoped) | — |
| Dead-letter exchange (DLX) | Dead-letter re-routing | Target queue channel |
Direct Reply-To (amq.rabbitmq.reply-to) | In-connector reply shortcut | Opaque per-channel address |
Vhost / | DefaultVhost segment | amqp.default.{queue} |
Vhost myvhost | Literal segment | amqp.myvhost.{queue} |
A message published over AMQP to queue orders in the default vhost lands in the KubeMQ Queue
channel amqp.default.orders. That channel is interoperable with gRPC/REST queue clients: an
AMQP producer can feed a native KubeMQ consumer on the same channel, and vice-versa. See
Channel mapping for the full grammar.
Vhost existence. There is no vhost-provisioning operation here — no rabbitmqctl add_vhost
equivalent — and connecting does not create one either. A vhost is materialized only by an
active exchange or queue declare. The pre-declared exchanges are a property of the vhost
name, so a passive declare of amq.direct succeeds against a vhost nothing has ever been
declared into, and the management API's exchange list agrees with the wire.
Canonical Client Example
The examples below use pika 1.x. Key API symbols:
pika.BlockingConnection, pika.URLParameters, pika.ConnectionParameters,
channel.queue_declare, channel.basic_publish, channel.basic_consume,
channel.basic_ack, channel.basic_get.
Work queue (publish + consume)
import pika # pika 1.x
params = pika.URLParameters("amqp://user:YOUR_JWT@kubemq.example.com:5672/")
conn = pika.BlockingConnection(params)
ch = conn.channel()
ch.queue_declare(queue="orders", durable=True)
# Publish
ch.basic_publish(
exchange="",
routing_key="orders",
body=b'{"order_id": "ORD-1"}',
)
# Consume (manual ack)
def on_message(ch, method, props, body):
print("received:", body)
ch.basic_ack(delivery_tag=method.delivery_tag)
ch.basic_qos(prefetch_count=10)
ch.basic_consume(queue="orders", on_message_callback=on_message)
ch.start_consuming()Topic exchange routing
ch.exchange_declare(exchange="logs", exchange_type="topic")
ch.queue_declare(queue="errors", durable=True)
ch.queue_bind(exchange="logs", queue="errors", routing_key="*.error")
ch.basic_publish(exchange="logs", routing_key="app.error", body=b"boom") # routed
ch.basic_publish(exchange="logs", routing_key="app.info", body=b"fyi") # not routedPublisher confirms
ch.confirm_delivery()
ch.basic_publish(
exchange="",
routing_key="orders",
body=b"payload",
mandatory=True,
)
# pika raises UnroutableError on basic.return when mandatory=True and no routeDead-letter exchange
ch.exchange_declare(exchange="dlx", exchange_type="fanout")
ch.queue_declare(queue="dead", durable=True)
ch.queue_bind(exchange="dlx", queue="dead", routing_key="")
ch.queue_declare(
queue="work",
durable=True,
arguments={"x-dead-letter-exchange": "dlx"},
)
# A consumer that basic_nack(..., requeue=False) on "work" dead-letters to "dead"
# with RabbitMQ-exact x-death headers appended.RPC (Direct Reply-To)
import uuid
reply_queue = "amq.rabbitmq.reply-to"
corr_id = str(uuid.uuid4())
# Responder (subscribe to the work queue and reply)
def handle_request(ch, method, props, body):
response = process(body)
ch.basic_publish(
exchange="",
routing_key=props.reply_to,
properties=pika.BasicProperties(correlation_id=props.correlation_id),
body=response,
)
ch.basic_ack(delivery_tag=method.delivery_tag)
# Requester
ch.basic_consume(queue=reply_queue, on_message_callback=on_reply, auto_ack=True)
ch.basic_publish(
exchange="",
routing_key="rpc_queue",
properties=pika.BasicProperties(
reply_to=reply_queue,
correlation_id=corr_id,
),
body=b"hello",
)One-shot pull (basic.get)
method, props, body = ch.basic_get(queue="orders", auto_ack=False)
if method:
print("got:", body)
ch.basic_ack(delivery_tag=method.delivery_tag)
else:
print("queue empty") # answered within about 50 ms on an empty queue (deviation 4)Security
Authentication — the credential store is consulted first, then the platform flag.
Connectors.Amqp.Enable defaults to true, so port 5672 is open on a stock server. Its
authentication posture is decided by Connectors.Amqp.Credentials first and only then by
Authentication.Enable:
- Credential store configured: both the username and the password are checked against
Connectors.Amqp.Credentials, and the platform-token path is not consulted for that connection —Authentication.Enabledoes not change this. - Authentication disabled, no credential store (default): the listener accepts any SASL
PLAIN/AMQPLAINcredentials. If KubeMQ is reachable from untrusted networks, either enable authentication, configure users, or firewall ports 5672/5671. - Authentication enabled, no credential store: the SASL
PLAIN/AMQPLAINpassword must be a valid KubeMQ JWT. The username is recorded for audit anduser-idchecks, never authenticated. - SASL
EXTERNAL: authenticates by verified TLS client certificate. It is off by default — setConnectors.Amqp.SslCertLogin = trueto enable it. It never uses the platform token, and when a credential store is configured it consults that store, so a certificate cannot authenticate as someone the store does not list. The identity comes fromConnectors.Amqp.SslCertLoginFrom. - SASL
ANONYMOUSis never offered, in any configuration, where RabbitMQ 4.3.4 advertises it by default (deviation 18).
Users, passwords and the connector-local credential store
A migrating RabbitMQ deployment usually has its own users and per-user permissions. The
connector can be given the same thing: Connectors.Amqp.Credentials, a list of
{Username, Password, Permissions, Tags} entries set in the config file or through the
settings API (there is no environment-variable form). The full recipe — the per-vhost
configure/write/read triple, tags, and how a rabbitmqctl set_permissions line maps onto one
entry — is in Users and permissions.
- Not configured (the default): nothing changes. The SASL password is still a KubeMQ platform token when authentication is on, and the username is still recorded but never checked.
- Configured: both the username and the password are checked, over
PLAINorAMQPLAIN, against this store — the platform token is not consulted at all for a connection authenticating this way. The username becomes the identity, encoded:orders_producerbecomesamqp-_user-orders__producer(the_is doubled).
Configuring users changes the authorization subject, and that is the single most important
consequence of turning this on. With no users configured, an authenticated AMQP client's
authorization subject is its token's client id. The moment Connectors.Amqp.Credentials gets
its first entry, every PLAIN/AMQPLAIN connection's subject becomes its encoded
username under the amqp-_user- prefix instead — so every existing authorization policy
written against the token's client id stops matching, silently, unless it is updated in the
same change. The server logs one WARN at startup naming exactly this consequence when users
are configured and Authorization.Enable = true; treat that line as a checklist item.
Credential checks are constant-time, and failures are indistinguishable. A wrong password,
an unknown username, and no credentials at all produce the identical connection.close(403)
/ ACCESS_REFUSED - authentication failed — measured against a real RabbitMQ 4.3.4 broker,
which returns the same text for all three. The distinguishing detail goes to the audit log only.
On an operator-managed Kubernetes cluster, fill this store from Secrets.
It has no environment variable for the operator's spec.envFromSecrets to supply, and the
settings page refuses every save on a server the operator manages. Plan such a cluster's
PLAIN/AMQPLAIN clients around one of these Secret-backed paths — or use certificate login
over EXTERNAL, which needs no store at all. Seed a user with spec.amqp.defaultUserSecretRef or load
your definitions export with spec.amqp.definitionsSecretRef; see
Users and permissions.
Certificate login (SASL EXTERNAL)
mTLS clients are identified by their certificate, not a username/password pair. SASL
EXTERNAL requires three things, and the first is the one that catches people:
Connectors.Amqp.SslCertLogin = true— it defaults to false, and it is deliberately not inferred from the server's TLS settings (TlsPortmust also be set).- Mutual TLS on the server:
Security.Cert,Security.KeyandSecurity.Caall set. Without the CA the listener never asks for a client certificate; without the certificate and key there is no TLS listener at all. Either wayEXTERNALis never offered, the server still starts, and it logs oneWARNnaming the security mode it found. - The TLS handshake presented a certificate the listener verified.
RabbitMQ needs two opt-ins here too — enabling the rabbitmq_auth_mechanism_ssl plugin and
listing EXTERNAL in auth_mechanisms — so this is the same shape, expressed as one setting.
The identity is controlled by Connectors.Amqp.SslCertLoginFrom, RabbitMQ's own setting name:
distinguished_name(the default, and RabbitMQ's own default): the certificate's full RFC-2253 distinguished name, e.g.CN=svc,OU=eng,O=Acme.common_name: the certificate subject's common name alone. The subject must carry exactly one non-empty common name; a certificate with none, an empty one, or several is refused.subject_alternative_name, RabbitMQ's third value, is not supported — any other value is a configuration error at startup, not a silent fallback.
Token expiry and refresh
This is a behavior change if you are upgrading, and it can disconnect clients that work
today. A connection whose token expires used to be served indefinitely; it is now closed at
the token's exp with connection.close(320) and the reply text
CONNECTION_FORCED - credential expired — the same code and text RabbitMQ 4.3.4's OAuth 2
plugin sends. There is no grace period and no configuration switch.
If you issue short-lived tokens, use connection.update-secret to hand the connection a new
one before the old expires. On success the deadline moves to the new token's expiry. Six
things refuse a refresh, and every one of them closes the connection 530:
- The new secret does not authenticate:
NOT_ALLOWED - Secret update failed. - The new secret authenticates but names no subject (an empty client id): the same text.
- The new secret would change the connection's identity:
NOT_ALLOWED - New secret was refused by one of the backends— the same refusal RabbitMQ 4.3.4 makes. Reconnect instead of trying to become someone else on an existing connection. - The connection logged in with a client certificate (
EXTERNAL):Secret update failed— it holds no token to replace. - A connector-local credential store is configured:
Secret update failed— no connection holds a token then. Authentication.Enable = false:Secret update failed— there is no credential to refresh.
The last three are refused whatever the new secret says. A token with no exp claim sets
no deadline and its connection is never evicted for age.
Which client id a policy must name
Authorization. When platform (Casbin) authorization is enabled, per-operation checks run
against the mapped channel name (amqp.{vhost}.{queue}), and exchanges are an authorization
resource of their own.
There are four subject derivations, and Authentication.Enable chooses between only two of
them. A policy file written for one derivation grants nothing under another, and nothing
warns you when the governing derivation changes — the policies simply stop matching.
| Derivation | Selected when | The subject a policy must name |
|---|---|---|
| Credential store | Connectors.Amqp.Credentials has at least one entry — before Authentication.Enable is consulted | amqp-_user- + the encoded username: ^amqp-_user-orders__producer$ for the user orders_producer |
SASL EXTERNAL | SslCertLogin = true and a verified client certificate — independent of Authentication.Enable; consults the credential store when one is configured | amqp-_cert- + the encoded certificate identity: ^amqp-_cert-CN_x3dsvc_x2cOU_x3deng$ for CN=svc,OU=eng |
| Platform token | no store and Authentication.Enable = true | the JWT's client id, verbatim and unprefixed: ^orders-producer$ |
| Anonymous | no store and Authentication.Enable = false | amqp- + 8 random characters per connection — no policy can name a particular client |
The amqp- forms are not amqp- + the name you configured. The name goes through an
encoder first: _ becomes __, and every other byte outside a-z A-Z 0-9 - becomes _x plus
two lowercase hex digits. So orders_producer → amqp-_user-orders__producer, svc.orders →
amqp-_user-svc_x2eorders, alice@corp.example → amqp-_user-alice_x40corp_x2eexample. A
policy naming the raw username matches nothing, and that user is denied every operation. Do not
verify the rule with a name like alice or orders-producer — those pass through unchanged and
teach the wrong rule.
Write one anchored literal per subject. Separators are escaped, because an unescaped . is a
wildcard:
[
{
"_rule": "AMQP producer, platform-token identity. Requires Authentication.Enable = true and NO Connectors.Amqp.Credentials.",
"Queues": true,
"ClientID": "^orders-producer$",
"Channel": "^amqp\\.default\\.orders$",
"Read": false,
"Write": true
},
{
"_rule": "AMQP consumer, same queue, read only.",
"Queues": true,
"ClientID": "^orders-worker$",
"Channel": "^amqp\\.default\\.orders$",
"Read": true,
"Write": false
}
]With Connectors.Amqp.Credentials configured, the same two rules name the encoded
usernames behind the amqp-_user- tag instead — orders_producer gains a second underscore
and orders-worker's body does not change, because - is in the identity charset and _ is
not:
[
{
"_rule": "AMQP producer, connector-local user 'orders_producer'.",
"Queues": true,
"ClientID": "^amqp-_user-orders__producer$",
"Channel": "^amqp\\.default\\.orders$",
"Read": false,
"Write": true
},
{
"_rule": "AMQP consumer, connector-local user 'orders-worker'.",
"Queues": true,
"ClientID": "^amqp-_user-orders-worker$",
"Channel": "^amqp\\.default\\.orders$",
"Read": true,
"Write": false
}
]Do not write amqp-.* as the client id — it is fail-open, and it is the shape a migrating
RabbitMQ operator reaches for first. Every configured username and every certificate
identity lands in the amqp- namespace, so that one rule matches every user in
Connectors.Amqp.Credentials at once; orders_producer and audit-reader become the same
principal, nothing is denied, and nothing is logged. With no credential store and
authentication off it is wider still — it matches every anonymous connection on a listener that
accepts any credentials. Write one anchored, encoded literal per user. If your AMQP clients
need different permissions from each other, configure Connectors.Amqp.Credentials or enable
authentication. With neither, AMQP authorization is all-or-nothing, and "all" means every
unauthenticated peer.
Publish is authorized on the exchange, not per destination queue — and then once more at the
shared queue layer. The connector checks Write on the exchange before routing; a denial
closes the channel with 403 naming the exchange and the publish is never acknowledged,
never returned as 312. Then every destination queue is authorized by the shared layer every
KubeMQ protocol goes through — a refusal there fails the whole publish as one outcome (a
basic.nack in confirm mode, a 541 channel close without confirms), never a partial delivery
and never a silent drop. This is stricter than RabbitMQ, which authorizes only the exchange.
Migrating a policy is therefore additive: keep every queue-scoped Write rule and add a
Write rule naming the exchange.
The platform layer has two verbs where RabbitMQ has three: queue.declare, queue.bind,
queue.unbind, queue.purge and queue.delete all check Write on the queue's channel, and
queue.bind/queue.unbind additionally need Read on the exchange. A policy that lets a client
declare a queue also lets it delete and publish to it. A connector-local user configured with
Permissions gets RabbitMQ's own three verbs, checked before this platform layer runs — see
Users and permissions. A platform-token
login, and a user with no Permissions entries, keep the two-verb model, so split producers and
topology owners onto different client ids if you relied on RabbitMQ's configure boundary.
Not authorized at all: any authenticated client can create server-named amq.gen-* queues
(queue.declare with an empty name — the generated name does not exist yet when the check
would run). Nothing meters or reclaims them beyond MaxConnections/ChannelMax.
TLS. Configure the server's Security block; the AMQPS listener on 5671 then serves the
same certificates as gRPC/REST. Set Connectors.Amqp.Port = 0 to force TLS-only AMQP.
Environment variables:
| Variable | Purpose |
|---|---|
CONNECTORS_AMQP_ENABLE | On by default; false closes the listener. An explicit true makes AMQP-specific configuration refusals a startup error instead of a warning |
CONNECTORS_AMQP_PORT | Plain TCP listener port (default 5672; 0 disables) |
CONNECTORS_AMQP_TLS_PORT | TLS listener port (default 5671; 0 disables) |
CONNECTORS_AMQP_SSL_CERT_LOGIN | true turns on certificate login over SASL EXTERNAL (default false) |
CONNECTORS_AMQP_SSL_CERT_LOGIN_FROM | distinguished_name (default) or common_name |
CONNECTORS_AMQP_MANAGEMENT_ENABLE / _PORT | The management HTTP API listener (default off, port 15672) |
See Authentication & security for JWT issuance and policy syntax, and the configuration reference for every field.
What Does NOT Migrate / Deviations
Never coming — and what to use instead
Each of these is a deliberate, permanent decision, not a gap awaiting work. Plan around them:
| RabbitMQ feature | Why it stays out | Use instead |
|---|---|---|
| Federation, shovel | Cross-broker topology plugins; KubeMQ moves data between clusters at the platform level | KubeMQ's own bridging, or definitions.json export/import to move topology |
Streams (x-queue-type: stream, x-stream-offset) | A different protocol and storage model | KubeMQ Events Store (persistent pub/sub with replay) |
Priority queues (x-max-priority) | The queue engine has no per-message priority in storage | Separate queues per priority class; consumer priorities (x-priority) do work |
| Consistent-hash exchange, deduplication plugin | Plugins with their own state | Partition by routing key at the publisher; deduplicate by message_id at the consumer |
Sender-selected distribution (CC/BCC headers) | Rarely used; ruled out until someone asks | Publish once per destination, or fan out through a fanout/topic exchange |
Lazy queues (x-queue-mode) | Moot — every queue is disk-backed | Nothing to do; the argument is accepted and ignored |
Firehose tracing (amq.rabbitmq.trace) | Nothing is published to it, though the exchange exists so tooling can bind | The audit log and the dashboard |
connection.blocked on memory/disk alarms | The queue engine has no such alarm; the notification fires only while the broker is not ready | Monitor the node's metrics |
| RabbitMQ-format Prometheus metrics | Different metric names | KubeMQ's own metrics endpoint |
| Requeue to head | Requeue is the storage engine's tail requeue (deviation 1) | Order-sensitive consumers should ack in order and avoid requeue |
Confirm multiple coalescing | Publishes complete in order one at a time, so there is never a batch; one ack frame per publish is spec-legal | Nothing to do; only a frame-count cost at very high rates |
Queue names with ; : * >, whitespace or a trailing . | They are KubeMQ channel separators, and escaping them would change the channel name every other protocol sees | Rename before migrating |
Topic permissions (set_topic_permissions) | Only the queue/exchange-scoped configure/write/read triple is enforced | Model the grant on the exchange and the bound queues |
Hard rejections (connection.close 540 not-implemented)
These AMQP methods are rejected immediately — your application will receive a protocol-level error and must stop using them. Each is exactly what RabbitMQ 4.3.4 answers too:
| Feature | Behavior |
|---|---|
basic.recover-async(requeue=false) | 540 — requeue=true is implemented (real requeue, no reply) |
basic.qos(prefetch-size != 0) | 540 — RabbitMQ refuses a byte-based prefetch rather than ignoring it |
immediate=true on basic.publish | 540 — same as RabbitMQ 4.x |
channel.flow(active=false) | 540 — RabbitMQ dropped client-driven flow control long ago; active=true replies flow-ok |
Inert arguments (accepted, ignored, logged)
The following declare arguments are accepted, stored in topology metadata, surfaced in the dashboard, and logged once per entity — they never alter behavior:
- Queue length limit in bytes (
x-max-length-bytes) — the queue-depth snapshot carries no byte total - Queue type (
x-queue-type), lazy mode (x-queue-mode) - Message priority (
x-max-priority)
x-max-length, x-overflow, x-message-ttl, x-expires, x-delivery-limit,
x-single-active-consumer, consumer x-priority and alternate-exchange are NOT inert. All
are enforced — and x-expires deletes an idle queue, x-max-length without x-overflow
is refused at declare time, and the dead-letter hop ceiling destroys a message that
crosses it. Read the deviations below before you upgrade a topology that sets any of them.
What works that older guides said did not
- Transactions (
tx.select/tx.commit/tx.rollback) work with RabbitMQ's measured semantics — publishes and settlements are held until commit, rollback discards publishes and leaves deliveries unacknowledged, and Spring AMQP'schannelTransacted=trueruns unchanged. One bound RabbitMQ does not have: an open transaction holds at most 65,536 operations and 128 MiB of message bodies; crossing it closes the channel406 transaction too large. - Exchange-to-exchange bindings (
exchange.bind/exchange.unbind) work on the wire and in a definitions import: one copy per queue however many paths reach it, cycles allowed,writeon the destination andreadon the source. - Alternate exchanges (
alternate-exchangeonexchange.declare): unroutable messages go to the alternate, chains and internal alternates included; amandatorypublish is returned only when the alternate cannot route it either. - Single active consumer (
x-single-active-consumer: true): only the first-registered consumer receives; the next in registration order takes over when it cancels or its connection drops, and a later consumer never pre-empts the active one whatever its priority. Held across the cluster by node presence. - Consumer priority (
basic.consumeargumentx-priority): while a higher-priority consumer has prefetch room it receives every message, lower-priority consumers get only the overflow, equal priorities share round-robin. Priorities are arbitrated per node. - The delayed-message plugin's exchange type (
x-delayed-message): declare it withx-delayed-typenaming the real routing type, publish with thex-delayheader, and the topology works unchanged —x-delayis honored on every exchange here, plugin or not. - Exchange
auto-delete: an auto-delete exchange is removed when its last binding goes, and one that was never bound stays, exactly as RabbitMQ 4.3.4 behaves. connection.update-secretrefreshes a platform token in place — see Token expiry and refresh.- An acknowledgement followed immediately by a disconnect is honored. Acknowledgements are applied at once, and the close waits at most a quarter of a second for the last one to land.
- Cluster-wide exclusivity and cross-node Direct Reply-To — see deviation 8.
Management HTTP API — what does not migrate
The management HTTP API covers a read-mostly
subset plus definitions.json import. Not implemented, at all: the management UI; every
write endpoint other than POST /api/definitions; /api/connections, /api/channels,
/api/nodes, /api/vhosts, /api/whoami, /api/health/checks/*; pagination and columns=;
rates, message_stats, memory and node fields; and federation/shovel management.
Behavioral deviations
These behaviors differ intentionally from RabbitMQ. Review each against your application before migrating. The numbering is stable and other pages refer to it.
- Requeue → tail. Requeued messages re-enter at the queue tail. RabbitMQ classic queues preserve near-head position.
- Expiry is eager on the delayed-retry shape, and lazy on everything else. A background
sweep dead-letters expired messages with nothing consuming, but only for a queue with a
queue-level
x-message-ttlgreater than 0, anx-dead-letter-exchange, and no consumer on any node — exactly the canonical wait-queue retry topology, which now fires on its own with no poller. Two shapes stay lazy and need a reader: a wait queue driven by a per-messageexpiration, and a queue with a TTL but no dead-letter exchange. The sweep visits at most 8 candidate queues per 30-second tick, so the worst case from deadline to dead-lettering isceil(candidates / 8) × 30 s— a 24-queue retry ladder can be up to 90 seconds late. The dead-letter exchange and its binding must exist before messages expire, or every swept message is dropped. On a dead-letter-bearing queue the deadline binds AMQP readers only: a gRPC, REST or dashboard reader of the same channel can still receive an expired message (stale, never lost). - TTL clamped and rounded up — on queues without a dead-letter exchange. A per-message TTL
is clamped to
Queue.MaxExpirationSeconds(default 12h) and rounded up to the next second. On a dead-letter-bearing queue the expiration is honored exactly as sent, to the millisecond.x-delayis clamped toMaxDelaySecondseverywhere. basic.getempty latency.basic.geton an empty queue answers within about 50 ms (a short broker poll, so the answer is authoritative). RabbitMQ answers in about 1 ms.MaxReceiveCountdrop. Messages redelivered beyond the brokerMaxReceiveCount(default 1024) are dropped or broker-rerouted. RabbitMQ redelivers forever. A rerouted message arrives on its dead-letter queue with the per-message TTL cleared.- DLX triggers —
rejected,expiredandmaxlen. RabbitMQ'sdelivery_limitreason string does not exist as its own value: a message that exceedsx-delivery-limitis dead-lettered with reasonrejected.basic.getgives up sweeping after 8 expired messages and answersget-emptyeven when a live message sits behind them; nothing is lost, the next call continues. - Inert queue arguments. Priority,
x-max-length-bytesand queue-type arguments are inert (see above).x-max-lengthandx-overfloware not in this group — see deviation 13. - Cluster-wide exclusivity, by gossip. Exclusive consumers, exclusive queue names and single-active queues are held across the cluster, and Direct Reply-To replies are forwarded to the requester's node. Two nodes that grant the same thing inside one gossip delay settle it deterministically (the earlier claim stands), and a node silent for 20 seconds is declared dead and its claims freed. Two residues: single-active takeover order is by node presence, and consumer priorities stay per node.
- Reserved vhost name and charset. The literal vhost name
default(the configuredDefaultVhost) is reserved — reach it via/. Queue and vhost names must not contain; : * >, whitespace, or a trailing.; vhost names may not contain.anywhere (402), because the channel name joins vhost and queue with a dot. Queue names may. - Transient non-exclusive queues are accepted here and refused by RabbitMQ 4.x. A
queue.declarewithdurable=falseandexclusive=falsesucceeds on KubeMQ; modern RabbitMQ rejects it with541by default. This runs the opposite direction from every other entry: topology first written against KubeMQ can fail if later pointed at RabbitMQ 4.x. - Partial-routing nack. Publisher-confirm
basic.nackon partial routing failure does not roll back queues that already accepted the message. A publisher retry may duplicate. queue.deletewithif-emptycan be refused when the count is unavailable. The count comes from the broker's monitoring endpoint; when it is unreachable, or while the connector still holds unacknowledged deliveries for the queue, the delete is refused406 PRECONDITION_FAILEDrather than assumed empty. A teardown script must be prepared for a 406 and retry. The unacknowledged-deliveries arm needs no outage at all — one ordinary manual-ack consumer with messages in flight is enough.x-max-lengthwith nox-overflowis refused at declare,406. This is the single most likely thing to surprise a migrating application. RabbitMQ's default overflow mode,drop-head, silently discards the oldest ready message — positional deletion the connector has no primitive for — so both the absent form and an explicitdrop-headare refused, and the error text names the substitute:x-overflow: reject-publish(abasic.nackin confirm mode, no signal outside it) orreject-publish-dlx(dead-lettered with reasonmaxlen). The limit itself is best-effort: the depth is a ~2-second cache, concurrent publishers race on one snapshot, a failed depth read admits the publish, and the count excludes delivered-but-unacknowledged messages. A fan-out publish is refused atomically when any one target is at its limit — sibling queues with no limit lose the message too.- Queue-level
x-message-ttlis enforced, with two divergences:x-message-ttl: 0means "no TTL" here (RabbitMQ: expire-on-arrival), and on a live declare anx-message-ttlaboveQueue.MaxExpirationSecondson a queue with no dead-letter exchange is refused,406, not clamped. x-expiresdeletes the queue — read this before upgrading if you have ever set it. A queue with no consumer, no redeclare and nobasic.getfor its window is deleted, and every message on it destroyed. Publishing does not count as use — that is RabbitMQ's own definition. A nightly-batch queue drained once a day by a job that attaches a consumer only then is idle for the other twenty-three hours. Audit every queue that carriesx-expiresand remove it from any queue whose reader is periodic rather than continuous — including queues read by a non-AMQP client. A connected non-AMQP consumer does block eviction, but a reader that only peeks does not, and a fully drained channel is invisible to the check. The sweep runs every 30 seconds.x-expiresis node-local on its clock, in a cluster. A consumer on any node protects the queue everywhere, but the activity clock is not shared: a queue kept alive only bybasic.gettraffic on one node can be expired by another, and the delete propagates. Keep a consumer attached rather than polling.x-delivery-limitis enforced, writing the per-queue poison cap: a limit ofNpermitsNdeliveries and the(N+1)th diverts the message to the dead-letter exchange. The declare is refused406unless the queue'sx-dead-letter-exchangeresolves to exactly one queue — no exchange, a missing one, or a fan-out are all refused, because otherwise the message would be silently destroyed at the limit.0is refused, and a value aboveQueue.MaxReceiveCountis refused naming the ceiling.- No SASL
ANONYMOUSmechanism. RabbitMQ 4.3.4 advertisesAMQPLAIN ANONYMOUS PLAINby default; KubeMQ advertisesPLAIN AMQPLAIN, addsEXTERNALonly for certificate login, and never advertisesANONYMOUS. A client that selects it is refusedconnection.close(503). A client that used it should sendPLAINinstead. The AMQP 1.0 connector on the same ports differs — it offersANONYMOUSwhenever authentication is disabled. - Every message is persisted, regardless of
delivery_mode. Every AMQP queue is backed by a disk-backed KubeMQ Queue channel, sodelivery_mode=1(transient) is stored exactly likedelivery_mode=2. Nothing is lost by this; a workload that relied on transient messages being cheaper to drop will see them survive a restart instead. - The dead-letter hop ceiling destroys the message.
DeadLetterMaxHops(default 16) caps thex-deathcount per queue and reason, and on a work-queue ↔ wait-queue retry ladder it is a cap on retry cycles: the 17th cycle drops the message, logged but gone, rather than dead-lettering it onward. Count your cycles and set it above them.
See Capabilities and Error codes for the exhaustive per-method error tables and wire-contract detail.
Moving topology with definitions export/import
RabbitMQ's definitions.json export/import is a common way to move a whole vhost's topology —
queues, exchanges, bindings, users and permissions — between brokers. KubeMQ implements the
import half against the same declare/bind code path an AMQP client's own
queue.declare/exchange.declare/queue.bind takes, plus an export in KubeMQ's own shapes.
The endpoint reference is Management HTTP API;
this is the migration recipe.
-
Export from RabbitMQ:
rabbitmqadmin --host rabbitmq.example.com --port 15672 export definitions.json -
Add the users the export names to KubeMQ's own configuration. Users and permissions are configuration, not runtime state —
POST /api/definitionsnever writesConnectors.Amqp.Credentials, andPasswordis a required field there. Add each user the export lists, with the sameTagsandPermissionsand aPasswordyou choose (RabbitMQ's export carries only apassword_hash). If the file carries apassword_hash, the import uses it to verify the configured password matches — not to set it — and reports a mismatch. -
Enable KubeMQ's management listener:
[Connectors.Amqp.Management] Enable = true Port = 15672 -
Import.
rabbitmqadminv2, or a plaincurl, both work — the listener speaks HTTP Basic like RabbitMQ's own:curl -u admin:adminpass -X POST \ -H 'Content-Type: application/json' \ --data-binary @definitions.json \ http://kubemq.example.com:15672/api/definitions -
Read the report. A clean import answers
204 No Content. Anything KubeMQ could not apply — policies, runtime parameters, a user not yet added toConnectors.Amqp.Credentials, an unsupported exchange type — answers400with areportarray naming every such entry, its kind, its vhost and why. Every entry is still attempted;applied/unappliedcount outcomes, so everything else in the file that could be applied already was. -
Fix what the report names, and re-import. Import is idempotent — a queue, exchange or binding already present from the first pass is a no-op on the second — so repeating steps 4–6 until the reply is
204is the normal way to converge.
Verification Smoke Test
For the basic round trip, see the quick start; this section adds deeper checks.
-
Confirm the listener is up. The connector is on by default, so start KubeMQ and confirm the log line
started insecure amqp listener(plain TCP) orstarted secure amqp listener(TLS) appears. If port 5672 is already taken on the host, the server logs an error and runs without the connector — free the port or move it withCONNECTORS_AMQP_PORT. -
Confirm cross-protocol visibility (optional). Query the KubeMQ internal API to verify the queue channel exists — the response should include
smoke-testin the queues list:curl http://localhost:8080/api/amqp/topology
See Also
Migrate from another broker
The ecosystem→connector map, the cross-protocol matrix, and guides for every supported broker.
Users and permissions
Connector-local users, the per-vhost configure/write/read triple, tags, and what a rabbitmqctl line maps onto.
Management HTTP API
The RabbitMQ-compatible subset on port 15672 — endpoints, tags, and definitions export/import.
Getting Started
Connect, declare a queue, and send and receive over AMQP 0-9-1 in a few steps.
Capabilities
Exactly what is supported, what is inert, what is enforced, and the deviations behind these behaviors.
Configuration reference
The complete Connectors.Amqp.* settings, ports, credentials, certificate login, and the management listener.
Was this page helpful?