Skip to content

Explanation

How RabbitMQ 4.x works inside: the Erlang/OTP process model, protocol handlers, exchanges and bindings, the three queue types (classic, quorum, stream), the Khepri metadata store that replaced Mnesia, federation and shovel, and the security model. Look-up tables (ports, defaults, matrices) live in Reference; tasks live in How-to Guides.

The 4.x shift in one paragraph

Between 4.0 (2024-09) and 4.3 (2026-04) RabbitMQ turned every replicated piece of state into a Raft-style log: classic queue mirroring was removed in 4.0, quorum queues and streams became the only replicated data types, Khepri became the default metadata store in 4.2 and the only one in 4.3, and AMQP 1.0 became a core protocol next to AMQP 0-9-1. The practical result is one failure model for the whole cluster: a majority of nodes must be up.

Component Overview

The diagram shows one node: protocol listeners hand messages to channel/session processes, which route through exchanges into queue processes; Khepri holds the topology that routing reads.

flowchart TB
    Client["Clients: AMQP 0-9-1, AMQP 1.0, MQTT, STOMP, Stream protocol"]
    subgraph Node["RabbitMQ node - Erlang VM (BEAM)"]
        subgraph Listeners["Protocol listeners"]
            AMQP["amqp091 / amqp10 reader<br/>(5672, 5671)"]
            MQTT["rabbitmq_mqtt<br/>(1883, 8883)"]
            STOMP["rabbitmq_stomp<br/>(61613)"]
            StreamProto["rabbitmq_stream<br/>(5552)"]
        end
        Chan["Channel / session processes<br/>(prefetch, confirms, acks)"]
        subgraph Routing["Exchange routing"]
            Direct["direct"]
            Topic["topic"]
            Fanout["fanout"]
            Headers["headers"]
            Plugins["x-consistent-hash, x-local-random,<br/>x-modulus-hash"]
        end
        subgraph Queues["Queue types"]
            CQ["Classic queue<br/>(CQv2, single replica)"]
            QQ["Quorum queue<br/>(Ra / Raft log + WAL)"]
            ST["Stream<br/>(Osiris segment log)"]
        end
        Khepri["Khepri metadata store<br/>(Ra / Raft)"]
        Mgmt["rabbitmq_management<br/>(15672)"]
        Prom["rabbitmq_prometheus<br/>(15692)"]
    end
    Client --> AMQP
    Client --> MQTT
    Client --> STOMP
    Client --> StreamProto
    AMQP --> Chan
    MQTT --> Chan
    STOMP --> Chan
    Chan --> Direct
    Chan --> Topic
    Chan --> Fanout
    Chan --> Headers
    Chan --> Plugins
    Direct --> QQ
    Topic --> QQ
    Topic --> ST
    Fanout --> CQ
    Headers --> QQ
    Plugins --> QQ
    StreamProto --> ST
    Khepri -.->|"vhosts, users, exchanges,<br/>bindings, policies"| Routing
    Khepri -.->|"queue records"| Queues

Components

Component Role
Erlang VM (BEAM) Lightweight processes with isolated heaps; RabbitMQ runs one process per connection, per channel/session and per queue replica.
Connection One Erlang process per client TCP connection. Handles heartbeats and multiplexes channels (AMQP 0-9-1) or sessions (AMQP 1.0).
Channel / session Cheaper than a TCP connection. Holds prefetch (QoS), unacked deliveries, publisher-confirm state. In 4.1+ quorum queue log reads are offloaded to channels, spreading CPU use.
Exchange Stateless routing element; has no message storage. Type chosen at declaration.
Binding Edge from exchange to queue (or exchange to exchange) with a routing key and optional arguments.
Classic queue Single-replica queue on one node (CQv2 storage). Still developed; used for transient, exclusive and RPC reply queues.
Quorum queue Raft-replicated FIFO queue built on the Ra library. The default choice when replication is needed.
Stream Append-only replicated log read by offset, built on the Osiris library.
Khepri Raft-based tree database (built on Ra) holding vhosts, users, permissions, topology, runtime parameters and policies. The only metadata store since 4.3.
Plugins Protocols (MQTT, STOMP, Stream, Web-*), inter-cluster links (Federation, Shovel), auth backends (OAuth 2.0, LDAP, HTTP), observability (Prometheus, Management, Tracing), peer discovery (K8s, AWS, Consul, etcd).

AMQP 0-9-1 Routing Model

The publisher never addresses a queue directly: it publishes to an exchange with a routing key, and bindings decide which queues receive a copy.

sequenceDiagram
    participant P as Producer
    participant C as Channel
    participant E as Exchange
    participant Q as Queue
    participant Cn as Consumer
    P->>C: basic.publish exchange=X routing_key=k
    C->>E: route(k)
    E->>Q: deliver to queues whose binding matches k
    Q->>Q: enqueue (persist if durable)
    Q-->>P: basic.ack (only if confirms enabled)
    Q->>Cn: basic.deliver (push, up to prefetch)
    Cn->>Q: basic.ack delivery-tag=N

The core exchange types are direct, topic, fanout and headers; 4.0 added x-local-random and 4.3 moved x-modulus-hash into core. See Reference: Exchange Types for the full list.

AMQP 1.0 (core since 4.0) uses the same exchanges and queues but addresses them with a v2 address format (/exchanges/:name/:key, /queues/:name); AMQP 1.0 clients can also declare topology. In 4.2 Direct Reply-To works for AMQP 1.0 and across protocols, and message interceptors can inspect or annotate messages on AMQP 0-9-1, AMQP 1.0 and MQTT.

Queue Types Deep Dive

The choice of queue type is the main design decision; this flowchart summarizes the official guidance.

flowchart TD
    Start["New queue"] --> Temp{"Temporary, exclusive<br/>or RPC reply?"}
    Temp -->|yes| Classic["Classic queue<br/>(non-replicated)"]
    Temp -->|no| Replay{"Consumers need replay,<br/>large fan-out or<br/>non-destructive reads?"}
    Replay -->|yes| Scale{"One stream near its<br/>throughput limit?"}
    Scale -->|no| Stream["Stream"]
    Scale -->|yes| Super["Super stream<br/>(partitioned)"]
    Replay -->|no| Durable{"Must survive<br/>node loss?"}
    Durable -->|yes| Quorum["Quorum queue"]
    Durable -->|no| Classic

Quorum Queue

A publish is confirmed only after a majority of members have written and fsynced the entry.

sequenceDiagram
    participant P as Publisher
    participant L as QQ leader
    participant F1 as Follower 1
    participant F2 as Follower 2
    participant C as Consumer
    P->>L: publish (confirm mode)
    L->>L: append to Ra log, WAL fsync
    L->>F1: append_entries
    L->>F2: append_entries
    F1-->>L: written and fsynced
    Note right of L: majority (2 of 3) reached
    L-->>P: basic.ack (confirm)
    L->>C: deliver
    C->>L: basic.ack
    L->>F1: replicate settle command
    L->>F2: replicate settle command
  • Built on Ra, the RabbitMQ team's Raft library; 4.3 moved to Ra 3.x and the 8th version of the quorum queue state machine.
  • Snapshots and log truncation keep the log bounded; 4.0 added checkpoints for sub-linear recovery, 4.3 added recovery snapshots and smarter snapshot throttling.
  • Poison-message protection: since 4.0 the default delivery-limit is 20; messages redelivered that many times are dead-lettered or dropped. The limit exists because an endless fail-requeue loop stops Raft log compaction and can fill the disk.
  • 4.3 additions: strict priority queues, delayed retry with backoff, per-queue consumer timeout, consumer_disconnected_timeout (default 60 s) for partitioned consumers, and roughly halved per-message memory in many scenarios.
  • Membership: the initial group size is 3 by default (members never share a node). Continuous Membership Reconciliation (CMR) is an opt-in feature that grows queues back to a target size when nodes join; it does not replace explicit rabbitmq-queues add_member / delete_member when nodes are removed permanently.

Classic Queue

  • Single replica on one node. Mirroring (the old "HA" policy) was deprecated in 3.9-era releases and removed in 4.0; mirroring keys in policies have no effect after upgrade.
  • CQv2 storage (per-queue index plus shared message store). Since 3.12 classic queues always page to disk aggressively, so "lazy mode" no longer applies; 4.3 rejects x-queue-mode and x-queue-version=1.
  • Right fit for exclusive, auto-delete and RPC reply queues; since 4.3, non-durable non-exclusive queues are denied by default.

Stream

Streams replicate differently from quorum queues: the leader writes segments and replicas pull them over dedicated TCP ports, while a Ra-based stream coordinator manages membership and leader election.

flowchart LR
    Pub["Publisher<br/>(stream protocol 5552<br/>or AMQP)"] --> Leader
    subgraph N1["Node 1"]
        Leader["Stream leader<br/>(Osiris writer)"]
        Seg1[("Segment files<br/>+ index + Bloom filter")]
        Leader --> Seg1
    end
    subgraph N2["Node 2"]
        R1["Replica<br/>(Osiris reader)"]
    end
    subgraph N3["Node 3"]
        R2["Replica<br/>(Osiris reader)"]
    end
    Leader -->|"replication<br/>(TCP 6000-6500)"| R1
    Leader -->|"replication<br/>(TCP 6000-6500)"| R2
    Coord["Stream coordinator<br/>(Ra cluster)"] -.->|"leader election,<br/>membership"| Leader
    Seg1 --> ConsA["Consumer A<br/>offset=first"]
    Seg1 --> ConsB["Consumer B<br/>offset=next"]
  • Append-only segment files (default 500 MB each); consumers attach at first, last, next, an offset or a timestamp, and reading does not remove data. Retention is by max-age and/or max-length-bytes, evaluated per segment.
  • Publisher confirms are sent once a quorum of replicas has the data, but streams do not fsync; they rely on the OS page cache. Quorum queues are the safer choice when every confirmed message must survive a simultaneous power loss.
  • Replicas are managed explicitly (rabbitmq-streams add_replica / delete_replica); a new cluster node hosts no stream replicas until added.
  • Super streams (3.11+) partition a logical stream into several streams behind a direct exchange; with single active consumer they keep per-partition order.
  • Server-side filtering: Bloom-filter chunk filtering (3.13), AMQP 1.0 property filter expressions (4.1) and SQL filter expressions (4.2) let a consumer receive only a subset of a stream.
  • Streams can be used over AMQP 0-9-1 and AMQP 1.0 (x-queue-type: stream), but the dedicated binary Stream protocol on 5552 is what reaches the highest throughput.

Khepri (Metadata Store)

Khepri stores everything except messages: users and permissions, vhosts, exchanges, queues and bindings, runtime parameters and policies. It is a tree-structured database built on Ra, so metadata, quorum queues and streams share one consensus algorithm.

Trait Mnesia (removed in 4.3) Khepri
Replication Mnesia multi-master replication Raft (Ra)
Network partition Conflicting writes on both sides; resolved by pause_minority, pause_if_all_down or autoheal (or ignored) Only the majority side accepts metadata changes; no partition-handling strategy to choose
Availability rule Minority side could keep serving with divergent state A majority of nodes must be online for the cluster to be available
Topic routing ETS tables Khepri projections; 4.3 uses a trie in an ordered_set ETS table for faster topic routing with many bindings
CLI rabbitmqctl rabbitmqctl (same commands); force_reset deprecated since 4.1

Timeline: experimental in 3.13 (not upgradeable to 4.x in place), fully supported and opt-in via the khepri_db feature flag in 4.0, default for new clusters in 4.2, the only store in 4.3. A 4.3 node boots Mnesia-based data by migrating it to Khepri, but the team recommends enabling khepri_db on 4.2 before upgrading.

Changed failure behavior

With Mnesia a minority of nodes could keep accepting topology changes. With Khepri, declaring a queue, binding or user fails unless a majority of nodes is reachable. This is the intended trade-off: consistency over availability for metadata.

Federation Plugin

Federation moves messages between separate clusters (or vhosts) over AMQP links. A downstream broker subscribes to an upstream exchange or queue and republishes locally; links tolerate WAN latency and intermittent connectivity, which clustering does not.

flowchart LR
    subgraph US["Upstream cluster (us-east)"]
        UpEx["exchange: orders"]
        UpQ["internal federation queue"]
        UpEx --> UpQ
    end
    subgraph EU["Downstream cluster (eu-west)"]
        Link["federation link<br/>(AMQP 0-9-1 client)"]
        DownEx["exchange: orders<br/>(policy federation-upstream-set)"]
        DownQ["local quorum queue"]
        Link --> DownEx --> DownQ
    end
    UpQ -->|"AMQPS over WAN"| Link

Shovel Plugin

A shovel is a configured consumer-plus-publisher: it consumes from a source (AMQP 0-9-1 or AMQP 1.0) and republishes to a destination, which can be another broker. Typical uses are migrations, draining a backlog to another cluster, or bridging protocols. Since 4.2 a local shovel protocol moves messages inside one cluster over inter-node links instead of TCP client connections, with higher throughput and lower resource use.

Connection and Channel Lifecycle

An AMQP 0-9-1 connection authenticates once and then opens channels, which cycle between publishing and consuming until closed.

stateDiagram-v2
    [*] --> TCP: open socket
    TCP --> Authenticated: SASL (PLAIN, AMQPLAIN, EXTERNAL, ANONYMOUS)
    Authenticated --> Tuned: connection.tune (frame_max, channel_max, heartbeat)
    Tuned --> ChannelOpen: channel.open
    ChannelOpen --> Idle
    Idle --> Publishing: basic.publish
    Idle --> Consuming: basic.consume
    Publishing --> Idle: confirm or return
    Consuming --> Idle: basic.cancel
    ChannelOpen --> ChannelClosed: channel.close or channel error
    ChannelClosed --> Tuned
    Tuned --> [*]: connection.close

Since 4.1, the pre-authentication frame_max is 8192 bytes (to fit larger JWTs), which is why old Node.js amqplib versions with a 4096-byte default fail to connect.

Comparison Hooks

  • vs Kafka: RabbitMQ wins on routing flexibility, per-message acknowledgements and AMQP compatibility. Kafka wins on log-replay throughput, consumer-group semantics and the analytics ecosystem.
  • vs NATS: NATS wins on latency and footprint. RabbitMQ wins on routing primitives, dead-lettering and protocol breadth.
  • vs Pulsar: Pulsar offers tiered storage and multi-tenancy built in. RabbitMQ uses vhosts plus policies for tenant boundaries and has no tiered storage for streams.

Security Model

RabbitMQ layers transport security (TLS on every listener), authentication (SASL mechanism plus one or more auth backends), vhost-scoped authorization, and management-plane roles. Backends, permission semantics and user tags are tabulated in Reference: Access Control; setup steps are in How-to Guides.

Authentication and authorization flow

A connection is authenticated by a chain of backends (auth_backends.N), then every operation is checked against the vhost permissions of that user.

sequenceDiagram
    participant App as Client app
    participant Node as RabbitMQ node
    participant IdP as OAuth 2.0 IdP (JWKS)
    App->>Node: TLS handshake on 5671
    App->>Node: SASL PLAIN, password = JWT
    Node->>IdP: fetch JWKS (cached, via issuer discovery)
    IdP-->>Node: signing keys
    Node->>Node: verify signature, exp, aud = resource_server_id
    Node->>Node: map scopes to configure/write/read per vhost
    Node-->>App: connection.open-ok
    App->>Node: queue.declare orders
    Node->>Node: check configure permission on queue orders

Authorization model

Permissions are three regular expressions per user per vhost (configure, write, read), matched against resource names; topic permissions add routing-key patterns. Vhosts are the isolation unit: users, queues and policies never cross vhost boundaries, which is why the 2026 advisories about cross-vhost disclosure were treated as security bugs.

Encryption

  • In transit: TLS is available on every listener (AMQP, Stream, MQTT, STOMP, management, Prometheus). Inter-node Erlang distribution is plaintext unless nodes run with -proto_dist inet_tls.
  • At rest: RabbitMQ does not encrypt queue, stream or Khepri files. Use dm-crypt/LUKS or encrypted cloud volumes.
  • Erlang cookie: nodes and CLI tools authenticate to each other with a shared cookie. Anyone with the cookie and network access to 25672 can run arbitrary code on the node, so the cookie is effectively a root credential for the cluster.

Threat model

Threat Mitigation
Default guest account Only allowed from localhost by default (loopback_users.guest = true); delete it or keep the restriction
Anonymous logins (4.0+ ANONYMOUS mechanism) anonymous_login_user = none in production
Management UI on a public IP Bind 15672 to a private network; put an authenticating proxy in front
Erlang distribution takeover Private network for 4369/25672, strong cookie, inet_tls distribution
Privilege escalation through the administrator or policymaker tag Grant tags sparingly; several 2026 advisories involved monitoring/policymaker users
Plugin supply chain Only enable plugins you need; take community plugins from their official releases and pin versions
Man-in-the-middle on AMQP TLS with verify_peer; close 5672 in production
Vhost escape Tight permission regexes; one vhost per tenant
Federation/Shovel credential exposure Credentials sit in runtime parameters (URIs); restrict who can read them and rotate them
Replay of messages Idempotent consumers or app-level message IDs; stream producers can deduplicate by producer name + publishing ID
Slow or stuck consumers Prefetch limits, consumer_timeout (quorum queues in 4.3), delivery-limit, monitoring of consumer utilisation
OAuth token replay Short token lifetimes, audience per cluster; AMQP 1.0 connections can refresh tokens (4.1+)
MQTT/WebSocket floods Rate-limit at ingress; require auth on 1883/8883; watch the retained-message store (see 2026 advisory list)

Sources