Migrating from AMQP 1.0
Point a native AMQP 1.0 client at KubeMQ: quick start, address-prefix mapping, a runnable go-amqp test, RPC and the transaction deviations.
If you have an application that speaks native AMQP 1.0 (ISO/IEC 19464) — using a client such as Azure/go-amqp, AMQP.NET Lite, or Apache Qpid Proton — you can point it at KubeMQ's built-in AMQP 1.0 connector by changing only the endpoint and, where required, the SASL credentials. Address prefixes select the KubeMQ messaging pattern, so simple publish/consume needs no app-code rewrite. This is a drop-in migration at the endpoint / client level.
If you are migrating a JMS application (Java) instead, see the Migrating from JMS guide; for ActiveMQ applications routed by client type, see Migrating from ActiveMQ.
Overview
The AMQP 1.0 connector exposes KubeMQ's Queues, Events, Events Store, Commands, and Queries patterns over the native AMQP 1.0 wire protocol. It listens on the same ports as the AMQP 0-9-1 connector — 5672 (plain / SASL) and 5671 (TLS / mTLS) — because both protocols share a single listener that routes each connection by its protocol header. No separate firewall rule is needed beyond what the AMQP port already allows.
- Canonical client (this guide):
github.com/Azure/go-amqp, tested with v1.7.0; newer 1.x releases are expected to work. For .NET shops, AMQP.NET Lite is a direct alternative; the connection-string and address conventions are the same, but the snippets below targetgo-amqp. - Opt-in default: The connector is disabled by default. Set
CONNECTORS_AMQP10_ENABLE=true(the10stays attached toAMQP), orEnable = trueunder[Connectors.Amqp10]in TOML; on Kubernetes, setspec.amqp10.enabled: true. Port 5672 is already open on a stock server for AMQP 0-9-1 (on by default), so an AMQP 1.0 client against an unconfigured server connects and is then closed during protocol negotiation rather than refused.
For the full wire-protocol contract, see the AMQP connector capabilities 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 AMQP 1.0 connector on
The connector is off by default. This command turns it on and publishes port 5672 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 amqp10 --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
With authentication on, pass a KubeMQ token as the SASL PLAIN password.
conn, err := amqp.Dial(ctx, "amqp://old-broker:5672", nil)conn, err := amqp.Dial(ctx, "amqp://localhost:5672", nil)Send and receive one message
Create a Go module, save main.go in it, and run it.
mkdir amqp-smoke
cd amqp-smoke
go mod init smoke
go get github.com/Azure/go-amqp@v1.7.0package main
import (
"context"
"fmt"
"log"
"time"
amqp "github.com/Azure/go-amqp"
)
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
// Messaging authentication is off by default, so no SASL options are needed.
conn, err := amqp.Dial(ctx, "amqp://localhost:5672", nil)
if err != nil {
log.Fatalf("dial: %v", err)
}
defer conn.Close()
session, err := conn.NewSession(ctx, nil)
if err != nil {
log.Fatalf("session: %v", err)
}
sender, err := session.NewSender(ctx, "/queues/smoke-test", nil)
if err != nil {
log.Fatalf("sender: %v", err)
}
if err := sender.Send(ctx, amqp.NewMessage([]byte("hello kubemq")), nil); err != nil {
log.Fatalf("send: %v", err)
}
_ = sender.Close(ctx)
receiver, err := session.NewReceiver(ctx, "/queues/smoke-test", &amqp.ReceiverOptions{Credit: 10})
if err != nil {
log.Fatalf("receiver: %v", err)
}
msg, err := receiver.Receive(ctx, nil)
if err != nil {
log.Fatalf("receive: %v", err)
}
if err := receiver.AcceptMessage(ctx, msg); err != nil {
log.Fatalf("accept: %v", err)
}
fmt.Printf("received: %s\n", msg.GetData())
}go run .You should see:
received: hello kubemqTo 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.
amqp10:
enabled: trueIf kmq installed your cluster, apply the file with kmq:
kmq deploy update --installation messaging --kube-context YOUR_KUBE_CONTEXT --values cluster-values.yaml --deadline 10mReplace:
YOUR_KUBE_CONTEXT— your cluster's kubectl context;kubectl config get-contextslists them.
Check that data_retained is true. You should see:
{
"installation_id": "…",
"state": "…",
"target": "…",
"data_retained": true,
"next_action": "…"
}If Helm installed your cluster, apply the file with the command in Apply a settings change.
Either way, check that the connector's Service exists:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq get svc messaging-amqp10Check that the row lists 5672/TCP. You should see:
NAME TYPE CLUSTER-IP EXTERNAL-IP PORT(S) AGE
messaging-amqp10 ClusterIP … <none> 5672/TCP,… …Reach the connector from your machine
Keep this running in a second terminal:
kubectl --context YOUR_KUBE_CONTEXT -n kubemq port-forward svc/messaging-amqp10 5672:5672You 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 address at a time, and keep your AMQP 1.0 broker until you have verified the move.
- Install KubeMQ for real: Choose your path.
Compatibility Matrix
The cells below describe the AMQP 1.0 column of the cross-protocol migration matrix.
| Dimension | Support | Notes |
|---|---|---|
| Drop-in level | endpoint / client | Change the host in the AMQP URI; prefix addresses with the pattern. No app-code rewrite for simple publish/consume. |
| Point-to-point queues | ✅ | /queues/<ch> → KubeMQ Queues; competing consumers, ack/nack, visibility. |
| Pub/sub (non-durable) | ✅ | /events/<ch> → Events; fire-hose fan-out, sender-settled (at-most-once). |
| Durable / persistent subscriptions | ✅ | /events-store/<ch> → Events Store; backed by the persistence engine, resume from last acked position. |
| Request / reply (RPC) | ✅ | /commands/<ch> / /queries/<ch> → Commands / Queries; hand-rolled reply receiver (see snippet). |
| Ordering guarantee | ⚠️ node-local | Within one node; ordering is not cluster-wide. |
| Transactions | ❌ | AMQP coordinator / declare / discharge frames are not implemented. |
| Dead-letter / redrive | ❌ no client DLQ | No client-settable DLQ over this protocol. Poison messages that exceed MaxReceiveCount are silently dropped by the broker, not delivered to a dead-letter address.¹ |
| Selectors / filtering | ✅ selectors (SQL92 subset) | apache.org:selector-filter:string on Events / Events Store links; not supported on /queues/ links. |
| Auth model | PLAIN (JWT) / EXTERNAL | SASL PLAIN: password = KubeMQ JWT. SASL EXTERNAL: mTLS, cert CN = ClientID. |
| TLS / mTLS | ✅ 5671 | Active when the top-level Security block is configured. |
| Top unsupported | transactions; durable-unsub node-local | See What Does Not Migrate. |
¹ There is no client-settable DLQ over this protocol. The AMQP 1.0 connector never marks
published messages for dead-lettering, so a poison message that exceeds MaxReceiveCount is
silently dropped by the broker rather than delivered to any consumable dead-letter address.
For a genuine client-facing DLQ, use the RabbitMQ (DLX) or AWS (redrive) connector.
Connection / Endpoint Migration
Change only the host and, if required, the SASL credentials. The AMQP 1.0 port is shared with AMQP 0-9-1 on the same listener, so no firewall change is needed beyond what the AMQP port already allows.
| Before (existing broker) | After (KubeMQ) | |
|---|---|---|
| Plain | amqp://broker:5672 | amqp://kubemq-host:5672 |
| TLS | amqps://broker:5671 | amqps://kubemq-host:5671 |
| Auth | Broker-specific credentials | SASL PLAIN, password = KubeMQ JWT |
// go-amqp — amqp.Dial
conn, err := amqp.Dial(ctx, "amqp://kubemq-host:5672",
&amqp.ConnOptions{
SASLType: amqp.SASLTypePlain("svc-orders", "<kubemq-jwt>"),
})When Authentication.Enable = false on the server, the connector also accepts SASL ANONYMOUS
and bare AMQP headers (no SASL) — convenient for local development.
Concept & Destination Mapping
The address prefix of the AMQP link selects the KubeMQ messaging pattern. The leading /
is optional (queues/orders ≡ /queues/orders).
| Source concept | AMQP 1.0 address | KubeMQ pattern | Channel name |
|---|---|---|---|
| Queue / P2P | /queues/<name> | Queues | <name> |
| Topic / pub-sub | /events/<name> | Events | <name> |
| Durable topic / persistent sub | /events-store/<name> | Events Store | <name> |
| Command (fire-and-forget RPC) | /commands/<name> | Commands | <name> |
| Query (request-response RPC) | /queries/<name> | Queries | <name> |
| RPC reply token | /responses/<RequestID> | reply path | connection-scoped |
| Temporary / dynamic node | source.dynamic or target.dynamic | in-memory mailbox | node-local |
Selectors (Events / Events Store only): attach a filter under the
apache.org:selector-filter:string descriptor on the receiver link source. The SQL92 subset
supported includes comparisons, AND / OR / NOT, BETWEEN, IN, LIKE, IS NULL, and
parentheses, evaluated against application-properties (and standard JMS headers).
Bare addresses: when no prefix is present the connector resolves by JMS terminus capability
hint (queue → Queues, topic → Events) or falls back to DefaultPattern (default:
"queues").
Interop with AMQP 0-9-1: channels produced over the AMQP 0-9-1 connector use the prefix
amqp.<vhost>.<queue>; an AMQP 1.0 client reaches the same data at
/queues/amqp.<vhost>.<queue>.
See the address mapping reference for the full grammar and the longest-prefix rule.
From other AMQP 1.0 brokers (Solace / Azure Service Bus)
The connector speaks standard AMQP 1.0, so non-ActiveMQ AMQP 1.0 clients migrate the
same way — change the endpoint, then map destinations to <pattern>/<channel>.
- Solace PubSub+ — a Solace AMQP 1.0 sender/receiver targets a queue or topic by name.
Remap the Solace destination to
queues/<name>(persistent) orevents/<name>(direct). Solace exclusive/non-exclusive durable topic endpoints map toevents-store/<name>durable subscriptions. Solace selectors map onto the pub/sub selector (events/-only). - Azure Service Bus — Service Bus AMQP 1.0 entities (
queues/<q>,topics/<t>/subscriptions/<s>) remap toqueues/<channel>andevents-store/<channel>(durable). Azure SB sessions, scheduled/deferred delivery, dead-lettering, and transactions have no KubeMQ equivalent — drop those features (see Capabilities). Azure SB'samqps://+ SAS-token auth maps to KubeMQ SASL PLAIN with a JWT.
For any AMQP 1.0 broker, the discipline is identical: explicit <pattern>/<channel>
addresses, continuous credit for at-most-once patterns, symbolic amqp:* error
conditions (never numeric codes), and the deviations below.
Canonical Client Example
Client:
github.com/Azure/go-amqp(tested with v1.7.0) Symbols used:amqp.Dial,conn.NewSession,session.NewSender,session.NewReceiver,sender.Send,amqp.NewMessage,receiver.Receive,receiver.AcceptMessage, plusmsg.Properties.ReplyTo/msg.Properties.CorrelationIDfor RPC.
Publish and consume (Queues)
package main
import (
"context"
"fmt"
"log"
amqp "github.com/Azure/go-amqp"
)
func main() {
ctx := context.Background()
// amqp.Dial establishes the TCP connection and SASL handshake.
conn, err := amqp.Dial(ctx, "amqp://kubemq-host:5672",
&amqp.ConnOptions{
SASLType: amqp.SASLTypePlain("svc-orders", "<kubemq-jwt>"),
})
if err != nil {
log.Fatal(err)
}
defer conn.Close()
// conn.NewSession opens an AMQP session.
sess, err := conn.NewSession(ctx, nil)
if err != nil {
log.Fatal(err)
}
// --- Publish ---
// session.NewSender attaches a sender link to /queues/orders.
snd, err := sess.NewSender(ctx, "/queues/orders", nil)
if err != nil {
log.Fatal(err)
}
// sender.Send transfers one message; amqp.NewMessage wraps the body.
if err := snd.Send(ctx, amqp.NewMessage([]byte(`{"id":"1","item":"widget"}`)), nil); err != nil {
log.Fatal(err)
}
snd.Close(ctx)
// --- Consume ---
// session.NewReceiver attaches a competing-consumer receiver on the same queue.
rcv, err := sess.NewReceiver(ctx, "/queues/orders", &amqp.ReceiverOptions{
Credit: 10, // grant initial link credit
})
if err != nil {
log.Fatal(err)
}
// receiver.Receive blocks until a message arrives.
msg, err := rcv.Receive(ctx, nil)
if err != nil {
log.Fatal(err)
}
fmt.Printf("received: %s\n", msg.GetData())
// receiver.AcceptMessage settles the delivery (DISPOSITION accepted → broker AckRange).
if err := rcv.AcceptMessage(ctx, msg); err != nil {
log.Fatal(err)
}
rcv.Close(ctx)
}Events Store (durable subscription)
// session.NewReceiver on /events-store/<ch> → durable Events Store subscription.
// The broker resumes from the last acknowledged position on reconnect.
rcv, err := sess.NewReceiver(ctx, "/events-store/audit", &amqp.ReceiverOptions{
Credit: 64,
Durability: amqp.DurabilityUnsettledState, // terminus expiry-policy 'never'
})Selectors (Events / Events Store)
// Attach a SQL92 selector on an events receiver link source filter.
// Selector is evaluated in the connector before delivery.
rcv, err := sess.NewReceiver(ctx, "/events/orders",
&amqp.ReceiverOptions{
Credit: 32,
Filters: []amqp.LinkFilter{
amqp.NewSelectorFilter("priority > 5 AND region = 'EU'"),
},
})RPC — hand-rolled reply receiver
The AMQP 1.0 connector has no library-level request/reply helper. You must create a
dynamic reply receiver yourself, set msg.Properties.ReplyTo to its address, and match
responses by CorrelationID. The example below uses the go-amqp symbols.
// 1. Open a dynamic receiver to serve as the reply address.
// session.NewReceiver with DynamicAddress=true → connector allocates
// a temporary node and returns its address in the ATTACH reply.
replyRcv, err := sess.NewReceiver(ctx, "", &amqp.ReceiverOptions{
Credit: 1,
DynamicAddress: true,
})
if err != nil {
log.Fatal(err)
}
replyAddr := replyRcv.Address() // the connector-assigned dynamic node address
// 2. Attach a sender to the command channel.
snd, err := sess.NewSender(ctx, "/commands/status", nil)
if err != nil {
log.Fatal(err)
}
// 3. Build the request message.
// msg.Properties.ReplyTo tells the connector where to send the response.
// msg.Properties.CorrelationID allows matching the reply to the request.
req := amqp.NewMessage([]byte(`{"service":"inventory"}`))
req.Properties = &amqp.MessageProperties{
ReplyTo: &replyAddr,
CorrelationID: "req-001",
}
// sender.Send dispatches the request to the Commands channel.
if err := snd.Send(ctx, req, nil); err != nil {
log.Fatal(err)
}
// 4. receiver.Receive blocks for the reply; the connector routes it to replyAddr.
reply, err := replyRcv.Receive(ctx, nil)
if err != nil {
log.Fatal(err)
}
fmt.Printf("reply correlation=%v body=%s\n",
reply.Properties.CorrelationID, reply.GetData())
// receiver.AcceptMessage acknowledges the reply delivery.
replyRcv.AcceptMessage(ctx, reply)
replyRcv.Close(ctx)
snd.Close(ctx)Security
Authentication
| Mechanism | When | Credential |
|---|---|---|
SASL PLAIN | Always available | password = KubeMQ JWT; username is recorded for audit only |
SASL ANONYMOUS | Authentication.Enable = false | No credentials; ClientID derived from container-id |
SASL EXTERNAL | mTLS with verified client certificate | Certificate CN becomes the ClientID — no JWT needed |
TLS / mTLS
The TLS listener on port 5671 activates only when the top-level Security block is
configured. Point clients at amqps://kubemq-host:5671. For mutual TLS, configure the server
to request client certificates; the certificate CN then serves as the connection ClientID.
Authorization
With Authorization.Enable = true, the connection's ClientID is checked against the Casbin
policy per link:
- Sender link (client → KubeMQ):
Writeon the resolved channel, checked at ATTACH. - Receiver link (KubeMQ → client):
Readon the resolved channel, checked at ATTACH. - Anonymous-terminus sender: per-message
Writecheck againstproperties.to(1024-entry LRU cache, 60 s TTL). /responses/<RequestID>reply token: no policy check (connection-scoped).
Minimal TOML configuration
[Connectors.Amqp10]
Enable = true
Port = 5672 # shared with [Connectors.Amqp] (0-9-1) via the same listener
TlsPort = 5671 # active only when [Security] is configuredWhat Does NOT Migrate / Deviations
Not supported
| Feature | Detail |
|---|---|
| AMQP transactions | The coordinator, declare, discharge, transactional acquisition, and transactional retirement performatives are not implemented. JMS SESSION_TRANSACTED sessions do not work; use AUTO_ACKNOWLEDGE or CLIENT_ACKNOWLEDGE. XA / two-phase commit is not available. |
rcv-settle-mode = second | Two-phase receiver settlement is not supported (DETACH amqp:not-implemented). Use first (the go-amqp default). |
AmqpSequence body sections | Rejected (amqp:not-implemented). Use Data sections or AmqpValue. |
| AMQP-over-WebSocket | No WebSocket binding (RFC 7395). Use raw TCP on 5672 / TLS on 5671. |
| SASL SCRAM / GSSAPI / Azure CBS | Only PLAIN, ANONYMOUS, and EXTERNAL are offered. |
| AMQP management (node create/delete) | Use the KubeMQ REST API or dashboard instead. |
| Client-settable DLQ / redrive | No client-settable DLQ over this protocol. See footnote ¹ in the Compatibility Matrix for the behavior and alternative-connector options (RabbitMQ DLX / AWS redrive). |
Behavioral deviations
| Area | Behavior |
|---|---|
| Durable-subscription unsubscribe is node-local | Durable Events Store subscription identities are tracked per node. The same durable subscription attached on two cluster nodes causes a durable-subscription client-ID conflict. An unsubscribe() call or durable detach is local to the node that owns the registration; if the client reconnects to a different node the durable state may not transfer cleanly. Retry the attach on conflict. |
| Pub/sub is at-most-once (sender-settled fire-hose) | /events/ links: events are dropped when link credit is 0 (kubemq_amqp10_events_dropped_no_credit_total), not buffered. Grant credit continuously for a durable-enough consumer. |
| Events Store stalled-credit link detach | A bounded per-link buffer (MaxUnsettledPerLink) fronts the Events Store subscription. If the buffer fills while credit stays at 0, the link is detached (amqp:resource-limit-exceeded) and the buffered window is dropped. Affected positions are already acked, so a durable re-attach resumes after them. Size MaxUnsettledPerLink to match the consumer's expected burst. |
released increments receive count | The broker increments ReceiveCount on redelivery after a released / modified{delivery-failed=true} settlement. This counts toward MaxReceiveCount — a strict AMQP reading would not count a release as a delivery attempt. |
Selectors on /queues/ links | Rejected (amqp:not-implemented). Selectors work only on Events and Events Store receivers. |
| Dynamic (temporary) nodes are node-local | A temp reply node lives in memory on the owning node only. Direct cross-connection sends to another connection's temp node work only within the same node. RPC replies that travel through the broker path (/responses/<id>) are unaffected. |
header.priority does not schedule | The field round-trips as the amqp10.priority tag but drives no priority ordering. |
| Config hot-reload | Changing Connectors.Amqp10.* requires a server restart. |
Verification Smoke Test
For the basic round trip, see the quick start; this section adds deeper checks.
If it fails
- If
amqp.Dialfails, check that the AMQP 1.0 connector is on and port 5672 is reachable. A TCP connect that succeeds and is then closed during protocol negotiation means the port is open for AMQP 0-9-1 but AMQP 1.0 is not enabled. - A refused connection on 5672 means the port is not published or is blocked: check the container's port mapping or your port-forward.
- If
sender.Sendtimes out, verify the JWT in the SASL PLAIN password field is valid. - For Events / Events Store, swap the address prefix and confirm fan-out to multiple receivers.
- For RPC, run the command snippet and confirm a reply with matching
CorrelationIDis received within theDefaultRpcTimeoutSecondswindow (default: 30 s).
See Also
Migrate from another broker
Choose the right KubeMQ wire-protocol connector for your existing broker.
Migrating from JMS
Swap a JMS ConnectionFactory to Qpid JMS over the AMQP 1.0 connector.
Migrating from ActiveMQ
Route an ActiveMQ workload onto KubeMQ by client type.
AMQP capabilities
Exactly what is supported and what is rejected on the AMQP 1.0 wire.
Error conditions
The symbolic amqp:* error conditions and what triggers each.
Configuration reference
Canonical AMQP 1.0 connector settings and env-var table.
Was this page helpful?
Migrating from ActiveMQ
Move ActiveMQ clients to KubeMQ by client type: Java through Qpid JMS, STOMP and MQTT by endpoint. Quick start included; OpenWire is not supported.
Migrating from JMS
Swap your JMS ConnectionFactory to Apache Qpid JMS over KubeMQ's AMQP 1.0 connector: quick start, destination mapping, a runnable test, XA gaps.