🎯 LESSON OBJECTIVE__HTMLTAG_68___
- ✅ Understand RabbitMQ cluster architecture and messaging patterns
- ✅ Deploy RabbitMQ 3-node cluster using Cluster Operator
- ✅ Configure quorum queues for HA
- ✅ Setup TLS encryption and authentication
- ✅ Monitor RabbitMQ cluster with Prometheus
- ✅ Best practices: message durability, DLQ, flow control
PART 1: RABBITMQ CLUSTER ARCHITECTURE
1.1. Messaging Patterns
graph LR
subgraph P2P["1️⃣ Point-to-Point"]
PP["Producer"] --> Q1["Queue"] --> C1A["Consumer"]
Q1 --> C1B["Consumer"]
end
subgraph FANOUT["2️⃣ Pub/Sub Fanout"]
FP["Producer"] --> EX1["Exchange"]
EX1 --> FQ1["Queue1"] --> FC1["Consumer1"]
EX1 --> FQ2["Queue2"] --> FC2["Consumer2"]
end
subgraph ROUTING["3️⃣ Routing"]
RP["Producer"] -->|"routing.key"| EX2["Exchange"]
EX2 -->|"match"| RQ["Queue"] --> RC["Consumer"]
end
style P2P fill:#0f172a,stroke:#3b82f6,color:#e2e8f0
style FANOUT fill:#0f172a,stroke:#7c3aed,color:#e2e8f0
style ROUTING fill:#0f172a,stroke:#15803d,color:#e2e8f0
1.2. RabbitMQ Cluster Architecture
graph TD
SVC["🌐 K8s Service ClusterIP<br/>rabbitmq.messaging:5672"]
SVC --> R0
SVC --> R1
SVC --> R2
subgraph CLUSTER["🐰 RabbitMQ Quorum Cluster"]
R0["🟢 rabbit-0<br/>Leader<br/>Mnesia"]
R1["🔵 rabbit-1<br/>Follower<br/>Mnesia"]
R2["🔵 rabbit-2<br/>Follower<br/>Mnesia"]
R0 <-->|"Raft replication"| R1
R1 <-->|"Raft replication"| R2
end
R0 --> PVC0["💾 Ceph PVC 20Gi"]
R1 --> PVC1["💾 Ceph PVC 20Gi"]
R2 --> PVC2["💾 Ceph PVC 20Gi"]
style SVC fill:#7c3aed,stroke:#a78bfa,color:#e2e8f0
style CLUSTER fill:#0f172a,stroke:#3b82f6,color:#e2e8f0
style R0 fill:#15803d,stroke:#22c55e,color:#e2e8f0
style R1 fill:#1e3a5f,stroke:#3b82f6,color:#e2e8f0
style R2 fill:#1e3a5f,stroke:#3b82f6,color:#e2e8f0
style PVC0 fill:#1e293b,stroke:#475569,color:#e2e8f0
style PVC1 fill:#1e293b,stroke:#475569,color:#e2e8f0
style PVC2 fill:#1e293b,stroke:#475569,color:#e2e8f0
Quorum Queue: Raft consensus → data replicated across majority Classic Queue: Only on 1 node (mirrored = deprecated)
| Feature | Classic Queue | Quorum Queue | Stream |
|---|---|---|---|
| Replication | None (mirrored deprecated) | Raft-based (majority)_ | _Append-only log |
| Data Safety | Low | High | High |
| Performance | Highest | Good (slightly lower) | Best for fan-out |
| Use Case | Temp/non-critical | Business-critical | Event streaming |
| Ordering_ | FIFO | FIFO | FIFO per partition |
PART 2: INSTALL RABBITMQ CLUSTER OPERATOR
2.1. Install Operator
# Install RabbitMQ Cluster Operator: kubectl apply -f https://github.com/rabbitmq/cluster-operator/releases/latest/download/cluster-operator.ymlVerify:
kubectl get pods -n rabbitmq-system
NAME READY STATUS
rabbitmq-cluster-operator-7f7d8b8bb9-xxxxx 1/1 Running
CRDs installed:
kubectl get crds | grep rabbitmq
rabbitmqclusters.rabbitmq.com
2.2. Create Namespace for Messaging
kubectl create namespace messaging
kubectl label namespace messaging purpose=message-brokers
PART 3: DEPLOY RABBITMQ HA CLUSTER
3.1. RabbitmqCluster CRD
# rabbitmq-cluster.yaml: apiVersion: rabbitmq.com/v1beta1 kind: RabbitmqCluster metadata: name: production-rmq namespace: messaging spec: replicas: 3image: rabbitmq:3.13-management
resources: requests: cpu: "500m" memory: "1Gi" limits: cpu: "2" memory: "2Gi"
persistence: storageClassName: ceph-block storage: 20Gi
rabbitmq: additionalConfig: | # Cluster formation cluster_formation.peer_discovery_backend = rabbit_peer_discovery_k8s cluster_formation.k8s.host = kubernetes.default.svc.cluster.local cluster_formation.k8s.address_type = hostname cluster_formation.node_cleanup.interval = 10 cluster_formation.node_cleanup.only_log_warning = true cluster_partition_handling = pause_minority
# Quorum queues as default: default_queue_type = quorum # Memory management: vm_memory_high_watermark.relative = 0.7 vm_memory_high_watermark_paging_ratio = 0.8 disk_free_limit.relative = 1.5 # Connection limits: channel_max = 2047 heartbeat = 60 # Consumer timeout (prevent stuck consumers): consumer_timeout = 3600000 # Management plugin stats collection: collect_statistics_interval = 10000 advancedConfig: | [ {rabbit, [ {tcp_listen_options, [ {backlog, 128}, {nodelay, true}, {linger, {true, 0}}, {exit_on_close, false} ]} ]} ]. additionalPlugins: - rabbitmq_prometheus - rabbitmq_shovel - rabbitmq_shovel_managementaffinity: podAntiAffinity: requiredDuringSchedulingIgnoredDuringExecution: - labelSelector: matchLabels: app.kubernetes.io/name: production-rmq topologyKey: kubernetes.io/hostname
override: statefulSet: spec: template: spec: topologySpreadConstraints: - maxSkew: 1 topologyKey: kubernetes.io/hostname whenUnsatisfiable: DoNotSchedule labelSelector: matchLabels: app.kubernetes.io/name: production-rmq
# Deploy:
kubectl apply -f rabbitmq-cluster.yaml
# Watch pods:
kubectl -n messaging get pods -w
# production-rmq-server-0 1/1 Running 0 60s
# production-rmq-server-1 1/1 Running 0 90s
# production-rmq-server-2 1/1 Running 0 120s
# Verify cluster:
kubectl -n messaging exec production-rmq-server-0 -- rabbitmqctl cluster_status
# Services created:
kubectl -n messaging get svc
# production-rmq ClusterIP 10.96.x.x 5672,15672,15692
# production-rmq-nodes ClusterIP None 4369,25672
PART 4: QUORUM QUEUES AND POLICY
4.1. Create Quorum Queues
# Lấy credentials: kubectl -n messaging get secret production-rmq-default-user -o jsonpath='{.data.username}' | base64 -d kubectl -n messaging get secret production-rmq-default-user -o jsonpath='{.data.password}' | base64 -dPort-forward management UI:
kubectl -n messaging port-forward svc/production-rmq 15672:15672
Via CLI - tạo quorum queue:
kubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare queue name=orders.created
queue_type=quorum durable=true
arguments='{"x-quorum-initial-group-size": 3, "x-delivery-limit": 5}'Queue với Dead Letter Exchange:
kubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare exchange name=orders type=topic durable=truekubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare exchange name=orders.dlx type=fanout durable=truekubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare queue name=orders.dlq
queue_type=quorum durable=truekubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare binding source=orders.dlx
destination=orders.dlq destination_type=queue
kubectl -n messaging exec production-rmq-server-0 --
rabbitmqadmin declare queue name=orders.process
queue_type=quorum durable=true
arguments='{"x-dead-letter-exchange": "orders.dlx", "x-delivery-limit": 3}'
4.2. Virtual Hosts & Users
# Create vhost cho mỗi service: kubectl -n messaging exec production-rmq-server-0 -- \ rabbitmqctl add_vhost /orderskubectl -n messaging exec production-rmq-server-0 --
rabbitmqctl add_vhost /paymentsCreate service users:
kubectl -n messaging exec production-rmq-server-0 --
rabbitmqctl add_user order_service "$(openssl rand -base64 24)"kubectl -n messaging exec production-rmq-server-0 --
rabbitmqctl set_permissions -p /orders order_service "orders\.." "orders\.." "orders\..*"Giới hạn permissions (principle of least privilege):
configure: "orders\..*" → chỉ configure queues prefix orders.
write: "orders\..*" → chỉ publish vào exchanges prefix orders.
read: "orders\..*" → chỉ consume từ queues prefix orders.
PART 5: TLS ENCRYPTION
# Sử dụng cert-manager:
apiVersion: cert-manager.io/v1
kind: Certificate
metadata:
name: rabbitmq-tls
namespace: messaging
spec:
secretName: rabbitmq-tls-secret
issuerRef:
name: cluster-ca-issuer
kind: ClusterIssuer
commonName: production-rmq.messaging.svc
dnsNames:
- production-rmq.messaging.svc
- production-rmq.messaging.svc.cluster.local
- "*.production-rmq-nodes.messaging.svc.cluster.local"
---
# Update RabbitmqCluster:
apiVersion: rabbitmq.com/v1beta1
kind: RabbitmqCluster
metadata:
name: production-rmq
namespace: messaging
spec:
tls:
secretName: rabbitmq-tls-secret
caSecretName: rabbitmq-ca-secret
disableNonTLSListeners: true # ⚠️ Force TLS only
PART 6: MONITORING
# PodMonitor cho Prometheus:
apiVersion: monitoring.coreos.com/v1
kind: PodMonitor
metadata:
name: rabbitmq-monitor
namespace: messaging
spec:
selector:
matchLabels:
app.kubernetes.io/name: production-rmq
podMetricsEndpoints:
- port: prometheus
interval: 15s
path: /metrics
---
# Key metrics:
# rabbitmq_queue_messages_ready — Messages waiting
# rabbitmq_queue_messages_unacked — Messages in-flight
# rabbitmq_queue_consumers — Consumer count
# rabbitmq_connections — Total connections
# rabbitmq_channels — Total channels
# rabbitmq_process_resident_memory_bytes — RAM usage
# rabbitmq_disk_space_available_bytes — Disk available
# Alerting rules:
apiVersion: monitoring.coreos.com/v1
kind: PrometheusRule
metadata:
name: rabbitmq-alerts
namespace: messaging
spec:
groups:
- name: rabbitmq
rules:
- alert: RabbitMQQueueBacklog
expr: rabbitmq_queue_messages_ready > 10000
for: 5m
labels:
severity: warning
annotations:
summary: "Queue {{ $labels.queue }} has {{ $value }} messages"
- alert: RabbitMQNoConsumers
expr: rabbitmq_queue_consumers == 0
for: 10m
labels:
severity: critical
annotations:
summary: "Queue {{ $labels.queue }} has no consumers"
- alert: RabbitMQHighMemory
expr: rabbitmq_process_resident_memory_bytes / rabbitmq_resident_memory_limit_bytes > 0.8
for: 5m
labels:
severity: warning
- alert: RabbitMQNodeDown
expr: rabbitmq_identity_info < 3
for: 1m
labels:
severity: critical
PART 7: APPLICATION INTEGRATION
# Application sử dụng RabbitMQ:
apiVersion: apps/v1
kind: Deployment
metadata:
name: order-service
namespace: default
spec:
template:
spec:
containers:
- name: order-service
env:
- name: RABBITMQ_URL
value: "amqps://order_service:$(RABBITMQ_PASSWORD)@production-rmq.messaging:5671/orders"
- name: RABBITMQ_PASSWORD
valueFrom:
secretKeyRef:
name: order-rmq-secret
key: password
# Python producer (pika):
import pika
import json
connection = pika.BlockingConnection(
pika.URLParameters(os.environ['RABBITMQ_URL'])
)
channel = connection.channel()
# Publish with delivery confirmation:
channel.confirm_delivery()
message = {"order_id": "12345", "amount": 100.00}
channel.basic_publish(
exchange='orders',
routing_key='orders.created',
body=json.dumps(message),
properties=pika.BasicProperties(
delivery_mode=2, # Persistent message
content_type='application/json',
message_id=str(uuid.uuid4())
)
)
💡 KEY TAKEAWAYS
- Quorum queues: Raft-based replication, always used for production
- Cluster Operator: Declarative RabbitMQ on K8s, auto-clustering
- 3-node cluster: Tolerates 1 node failure (majority = 2)
- Dead Letter Queue: Handle poison messages, retry logic
- TLS + vhosts: Isolate services, encrypt in transit
- pause_minority: Prevent split-brain partition handling
🎯 EXERCISE
Exercise 1: RabbitMQ HA Lab
- Deploy 3-node RabbitMQ cluster
- Create quorum queue with DLQ
- Publish 10,000 messages, kill 1 node, verify no message loss
Exercise 2: Monitoring Setup__HTMLTAG_229___
- Configure PodMonitor
- Import RabbitMQ Grafana dashboard (ID: 10991)
- Create alert for queue backlog > 5000
📚 NEXT POST
In Lesson 22: Apache Kafka Cluster with Strimzi Operator, we will deploy Kafka for high-throughput event streaming.