KubeMQ
ConnectorsRabbitMQ (AMQP 0-9-1)Reference

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

ConnectorKubeMQ RabbitMQ connector (AMQP 0-9-1 RabbitMQ dialect)
Ports5672 (AMQP plain) / 5671 (AMQPS/TLS); 15672 for the optional management HTTP API
Canonical clientpika 1.x (Python)
Drop-in levelEndpoint-only — change the connection string; keep all client library code unchanged for the supported feature set
Enable defaultOn — 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.

Terminal
kmq deploy update --installation kubemq --connector amqp --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
amqp://user:password@rabbitmq.example.com:5672/
After
amqp://user:password@localhost:5672/

Send and receive one message

Install the client with pip install pika, then run each file with python3.

publish.py
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")
consume.py
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:

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

Terminal 2
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-amqp 5672:5672

Replace:

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

You should see:

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

Send 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.

DimensionRabbitMQ 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 modelSASL 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 unsupportedSee 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:

AMQP URI swap
# 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/myvhost

Key differences to be aware of at connection time:

AspectRabbitMQKubeMQ
SASL mechanismAMQPLAIN ANONYMOUS PLAIN advertised by default (measured on RabbitMQ 4.3.4); EXTERNAL once its plugin is enabledPLAIN and AMQPLAIN always; EXTERNAL only when Connectors.Amqp.SslCertLogin = true and the listener verified a client certificate. ANONYMOUS is never offered (deviation 18)
PasswordRabbitMQ user passwordNo 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
VhostMust be pre-createdNamespace-on-first-use — any charset-valid vhost is accepted; / maps to the configured DefaultVhost (default "default"). A vhost name may not contain . anywhere
Client identityConnection name (optional)Four derivations — see Which client id a policy must name. connection_name is display/log only and never reaches the identity
Token lifetimen/aA connection whose JWT expires is closed with connection.close(320); connection.update-secret refreshes it in place — see Token expiry and refresh
TLSPer-listener configActive when the server's Security block is configured; serves the same certificates as gRPC/REST
Queue storageRabbitMQ queuesKubeMQ 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 conceptKubeMQ patternChannel / address
QueueQueuesamqp.{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-routingTarget queue channel
Direct Reply-To (amq.rabbitmq.reply-to)In-connector reply shortcutOpaque per-channel address
Vhost /DefaultVhost segmentamqp.default.{queue}
Vhost myvhostLiteral segmentamqp.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 routed

Publisher 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 route

Dead-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.Enable does not change this.
  • Authentication disabled, no credential store (default): the listener accepts any SASL PLAIN/AMQPLAIN credentials. 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/AMQPLAIN password must be a valid KubeMQ JWT. The username is recorded for audit and user-id checks, never authenticated.
  • SASL EXTERNAL: authenticates by verified TLS client certificate. It is off by default — set Connectors.Amqp.SslCertLogin = true to 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 from Connectors.Amqp.SslCertLoginFrom.
  • SASL ANONYMOUS is 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 PLAIN or AMQPLAIN, against this store — the platform token is not consulted at all for a connection authenticating this way. The username becomes the identity, encoded: orders_producer becomes amqp-_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:

  1. Connectors.Amqp.SslCertLogin = true — it defaults to false, and it is deliberately not inferred from the server's TLS settings (TlsPort must also be set).
  2. Mutual TLS on the server: Security.Cert, Security.Key and Security.Ca all 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 way EXTERNAL is never offered, the server still starts, and it logs one WARN naming the security mode it found.
  3. 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.

DerivationSelected whenThe subject a policy must name
Credential storeConnectors.Amqp.Credentials has at least one entry — before Authentication.Enable is consultedamqp-_user- + the encoded username: ^amqp-_user-orders__producer$ for the user orders_producer
SASL EXTERNALSslCertLogin = true and a verified client certificate — independent of Authentication.Enable; consults the credential store when one is configuredamqp-_cert- + the encoded certificate identity: ^amqp-_cert-CN_x3dsvc_x2cOU_x3deng$ for CN=svc,OU=eng
Platform tokenno store and Authentication.Enable = truethe JWT's client id, verbatim and unprefixed: ^orders-producer$
Anonymousno store and Authentication.Enable = falseamqp- + 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:

authorization policy — platform-token identities
[
  {
    "_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:

authorization policy — connector-local users
[
  {
    "_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:

VariablePurpose
CONNECTORS_AMQP_ENABLEOn by default; false closes the listener. An explicit true makes AMQP-specific configuration refusals a startup error instead of a warning
CONNECTORS_AMQP_PORTPlain TCP listener port (default 5672; 0 disables)
CONNECTORS_AMQP_TLS_PORTTLS listener port (default 5671; 0 disables)
CONNECTORS_AMQP_SSL_CERT_LOGINtrue turns on certificate login over SASL EXTERNAL (default false)
CONNECTORS_AMQP_SSL_CERT_LOGIN_FROMdistinguished_name (default) or common_name
CONNECTORS_AMQP_MANAGEMENT_ENABLE / _PORTThe 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 featureWhy it stays outUse instead
Federation, shovelCross-broker topology plugins; KubeMQ moves data between clusters at the platform levelKubeMQ's own bridging, or definitions.json export/import to move topology
Streams (x-queue-type: stream, x-stream-offset)A different protocol and storage modelKubeMQ Events Store (persistent pub/sub with replay)
Priority queues (x-max-priority)The queue engine has no per-message priority in storageSeparate queues per priority class; consumer priorities (x-priority) do work
Consistent-hash exchange, deduplication pluginPlugins with their own statePartition 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 asksPublish once per destination, or fan out through a fanout/topic exchange
Lazy queues (x-queue-mode)Moot — every queue is disk-backedNothing 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 bindThe audit log and the dashboard
connection.blocked on memory/disk alarmsThe queue engine has no such alarm; the notification fires only while the broker is not readyMonitor the node's metrics
RabbitMQ-format Prometheus metricsDifferent metric namesKubeMQ's own metrics endpoint
Requeue to headRequeue is the storage engine's tail requeue (deviation 1)Order-sensitive consumers should ack in order and avoid requeue
Confirm multiple coalescingPublishes complete in order one at a time, so there is never a batch; one ack frame per publish is spec-legalNothing 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 seesRename before migrating
Topic permissions (set_topic_permissions)Only the queue/exchange-scoped configure/write/read triple is enforcedModel 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:

FeatureBehavior
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.publish540 — 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's channelTransacted=true runs 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 channel 406 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, write on the destination and read on the source.
  • Alternate exchanges (alternate-exchange on exchange.declare): unroutable messages go to the alternate, chains and internal alternates included; a mandatory publish 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.consume argument x-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 with x-delayed-type naming the real routing type, publish with the x-delay header, and the topology works unchanged — x-delay is 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-secret refreshes 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.

  1. Requeue → tail. Requeued messages re-enter at the queue tail. RabbitMQ classic queues preserve near-head position.
  2. 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-ttl greater than 0, an x-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-message expiration, 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 is ceil(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).
  3. 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-delay is clamped to MaxDelaySeconds everywhere.
  4. basic.get empty latency. basic.get on an empty queue answers within about 50 ms (a short broker poll, so the answer is authoritative). RabbitMQ answers in about 1 ms.
  5. MaxReceiveCount drop. Messages redelivered beyond the broker MaxReceiveCount (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.
  6. DLX triggers — rejected, expired and maxlen. RabbitMQ's delivery_limit reason string does not exist as its own value: a message that exceeds x-delivery-limit is dead-lettered with reason rejected. basic.get gives up sweeping after 8 expired messages and answers get-empty even when a live message sits behind them; nothing is lost, the next call continues.
  7. Inert queue arguments. Priority, x-max-length-bytes and queue-type arguments are inert (see above). x-max-length and x-overflow are not in this group — see deviation 13.
  8. 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.
  9. Reserved vhost name and charset. The literal vhost name default (the configured DefaultVhost) 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.
  10. Transient non-exclusive queues are accepted here and refused by RabbitMQ 4.x. A queue.declare with durable=false and exclusive=false succeeds on KubeMQ; modern RabbitMQ rejects it with 541 by default. This runs the opposite direction from every other entry: topology first written against KubeMQ can fail if later pointed at RabbitMQ 4.x.
  11. Partial-routing nack. Publisher-confirm basic.nack on partial routing failure does not roll back queues that already accepted the message. A publisher retry may duplicate.
  12. queue.delete with if-empty can 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 refused 406 PRECONDITION_FAILED rather 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.
  13. x-max-length with no x-overflow is 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 explicit drop-head are refused, and the error text names the substitute: x-overflow: reject-publish (a basic.nack in confirm mode, no signal outside it) or reject-publish-dlx (dead-lettered with reason maxlen). 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.
  14. Queue-level x-message-ttl is enforced, with two divergences: x-message-ttl: 0 means "no TTL" here (RabbitMQ: expire-on-arrival), and on a live declare an x-message-ttl above Queue.MaxExpirationSeconds on a queue with no dead-letter exchange is refused, 406, not clamped.
  15. x-expires deletes the queue — read this before upgrading if you have ever set it. A queue with no consumer, no redeclare and no basic.get for 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 carries x-expires and 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.
  16. x-expires is 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 by basic.get traffic on one node can be expired by another, and the delete propagates. Keep a consumer attached rather than polling.
  17. x-delivery-limit is enforced, writing the per-queue poison cap: a limit of N permits N deliveries and the (N+1)th diverts the message to the dead-letter exchange. The declare is refused 406 unless the queue's x-dead-letter-exchange resolves 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. 0 is refused, and a value above Queue.MaxReceiveCount is refused naming the ceiling.
  18. No SASL ANONYMOUS mechanism. RabbitMQ 4.3.4 advertises AMQPLAIN ANONYMOUS PLAIN by default; KubeMQ advertises PLAIN AMQPLAIN, adds EXTERNAL only for certificate login, and never advertises ANONYMOUS. A client that selects it is refused connection.close(503). A client that used it should send PLAIN instead. The AMQP 1.0 connector on the same ports differs — it offers ANONYMOUS whenever authentication is disabled.
  19. Every message is persisted, regardless of delivery_mode. Every AMQP queue is backed by a disk-backed KubeMQ Queue channel, so delivery_mode=1 (transient) is stored exactly like delivery_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.
  20. The dead-letter hop ceiling destroys the message. DeadLetterMaxHops (default 16) caps the x-death count 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.

  1. Export from RabbitMQ:

    rabbitmqadmin --host rabbitmq.example.com --port 15672 export definitions.json
  2. Add the users the export names to KubeMQ's own configuration. Users and permissions are configuration, not runtime state — POST /api/definitions never writes Connectors.Amqp.Credentials, and Password is a required field there. Add each user the export lists, with the same Tags and Permissions and a Password you choose (RabbitMQ's export carries only a password_hash). If the file carries a password_hash, the import uses it to verify the configured password matches — not to set it — and reports a mismatch.

  3. Enable KubeMQ's management listener:

    [Connectors.Amqp.Management]
      Enable = true
      Port = 15672
  4. Import. rabbitmqadmin v2, or a plain curl, 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
  5. Read the report. A clean import answers 204 No Content. Anything KubeMQ could not apply — policies, runtime parameters, a user not yet added to Connectors.Amqp.Credentials, an unsupported exchange type — answers 400 with a report array naming every such entry, its kind, its vhost and why. Every entry is still attempted; applied/unapplied count outcomes, so everything else in the file that could be applied already was.

  6. 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 204 is the normal way to converge.

Verification Smoke Test

For the basic round trip, see the quick start; this section adds deeper checks.

  1. 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) or started 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 with CONNECTORS_AMQP_PORT.

  2. Confirm cross-protocol visibility (optional). Query the KubeMQ internal API to verify the queue channel exists — the response should include smoke-test in the queues list:

    curl http://localhost:8080/api/amqp/topology

See Also

Was this page helpful?

On this page