概述
RabbitMQ 是一个开源的消息代理和队列服务器,用于通过普通协议在完全不同的应用之间共享数据。RabbitMQ 是使用 Erlang 语言开发的 AMQP (Advanced Message Queuing Protocol) 的实现。
核心特性
- 可靠性:支持消息持久化、确认机制、事务等
- 灵活的路由:支持多种交换机类型和路由模式
- 集群支持:支持集群部署和高可用性
- 多协议支持:支持 AMQP、MQTT、STOMP 等协议
- 管理界面:提供 Web 管理界面和命令行工具
- 插件系统:丰富的插件生态系统
应用场景
| 场景 | 描述 | 示例 |
|---|---|---|
| 异步处理 | 提高系统响应速度 | 用户注册后发送邮件 |
| 应用解耦 | 减少系统间依赖 | 订单系统 → 库存系统 |
| 流量削峰 | 平滑处理高并发 | 秒杀活动、双11购物 |
| 消息分发 | 一对多消息传递 | 系统通知、日志收集 |
💡 提示: RabbitMQ 特别适合需要高可靠性、复杂路由的企业级应用场景。
环境准备
系统要求
- Docker 20.10+
- 内存:至少 1GB RAM
- 存储:至少 5GB 可用空间
- 网络:确保端口 5672 和 15672 可访问
架构组件
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Producer │ │ Exchange │ │ Queue │
│ (生产者) │───▶│ (交换机) │───▶│ (队列) │
└─────────────┘ └─────────────┘ └─────────────┘
│ │
▼ ▼
┌─────────────┐ ┌─────────────┐
│ Binding │ │ Consumer │
│ (绑定) │ │ (消费者) │
└─────────────┘ └─────────────┘
Docker 安装部署
1. 获取 RabbitMQ 镜像
# 搜索 RabbitMQ 镜像
docker search rabbitmq
# 拉取带管理界面的镜像
docker pull rabbitmq:3.12-management
2. 单机部署
基本启动命令
# 基础启动(临时测试)
docker run -d \
--hostname my-rabbit \
--name rabbitmq \
-p 15672:15672 \
-p 5672:5672 \
rabbitmq:3.12-management
# 生产环境启动(推荐)
docker run -d \
--restart=always \
--name rabbitmq \
--hostname rabbitmq-node1 \
-p 5672:5672 \
-p 15672:15672 \
-p 25672:25672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=admin123 \
-e RABBITMQ_DEFAULT_VHOST=/ \
-v /etc/localtime:/etc/localtime \
-v rabbitmq_data:/var/lib/rabbitmq \
--log-driver json-file \
--log-opt max-size=100m \
--log-opt max-file=2 \
rabbitmq:3.12-management
3. 使用 Docker Compose 部署
创建 docker-compose.yml 文件:
version: '3.8'
services:
rabbitmq:
image: rabbitmq:3.12-management
container_name: rabbitmq
restart: unless-stopped
hostname: rabbitmq-node1
ports:
- "5672:5672" # AMQP 端口
- "15672:15672" # 管理界面端口
- "25672:25672" # 集群通信端口
environment:
# 默认用户和密码
RABBITMQ_DEFAULT_USER: admin
RABBITMQ_DEFAULT_PASS: admin123
RABBITMQ_DEFAULT_VHOST: /
# 启用管理插件
RABBITMQ_PLUGINS: rabbitmq_management
# 日志级别
RABBITMQ_LOG_LEVEL: info
volumes:
- rabbitmq_data:/var/lib/rabbitmq
- rabbitmq_logs:/var/log/rabbitmq
- /etc/localtime:/etc/localtime:ro
logging:
driver: json-file
options:
max-size: "100m"
max-file: "2"
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "ping"]
interval: 30s
timeout: 10s
retries: 3
start_period: 40s
volumes:
rabbitmq_data:
rabbitmq_logs:
启动服务:
# 启动 RabbitMQ
docker-compose up -d
# 查看日志
docker-compose logs -f rabbitmq
# 停止服务
docker-compose down
4. 集群部署
创建 docker-compose-cluster.yml:
version: '3.8'
services:
rabbitmq1:
image: rabbitmq:3.12-management
container_name: rabbitmq1
hostname: rabbitmq1
environment:
RABBITMQ_ERLANG_COOKIE: 'rabbitmq-cluster-cookie'
RABBITMQ_DEFAULT_USER: admin
RABBITMQ_DEFAULT_PASS: admin123
ports:
- "5672:5672"
- "15672:15672"
volumes:
- rabbitmq1_data:/var/lib/rabbitmq
networks:
- rabbitmq-cluster
rabbitmq2:
image: rabbitmq:3.12-management
container_name: rabbitmq2
hostname: rabbitmq2
environment:
RABBITMQ_ERLANG_COOKIE: 'rabbitmq-cluster-cookie'
ports:
- "5673:5672"
- "15673:15672"
volumes:
- rabbitmq2_data:/var/lib/rabbitmq
networks:
- rabbitmq-cluster
depends_on:
- rabbitmq1
rabbitmq3:
image: rabbitmq:3.12-management
container_name: rabbitmq3
hostname: rabbitmq3
environment:
RABBITMQ_ERLANG_COOKIE: 'rabbitmq-cluster-cookie'
ports:
- "5674:5672"
- "15674:15672"
volumes:
- rabbitmq3_data:/var/lib/rabbitmq
networks:
- rabbitmq-cluster
depends_on:
- rabbitmq1
volumes:
rabbitmq1_data:
rabbitmq2_data:
rabbitmq3_data:
networks:
rabbitmq-cluster:
driver: bridge
基本概念
核心组件
| 组件 | 描述 | 功能 |
|---|---|---|
| Producer | 生产者 | 发送消息的应用程序 |
| Consumer | 消费者 | 接收消息的应用程序 |
| Queue | 队列 | 存储消息的缓冲区 |
| Exchange | 交换机 | 接收消息并路由到队列 |
| Binding | 绑定 | 连接交换机和队列的规则 |
| Routing Key | 路由键 | 消息的路由标识 |
| Virtual Host | 虚拟主机 | 逻辑分组和权限控制 |
交换机类型
1. Direct Exchange (直连交换机)
Producer → [routing_key] → Exchange → Queue
完全匹配路由键
2. Topic Exchange (主题交换机)
Producer → [pattern] → Exchange → Queue
支持通配符匹配
3. Fanout Exchange (广播交换机)
Producer → Exchange → All Queues
广播到所有绑定的队列
4. Headers Exchange (头交换机)
Producer → [headers] → Exchange → Queue
基于消息头路由
核心功能
1. 用户和权限管理
创建用户
# 进入 RabbitMQ 容器
docker exec -it rabbitmq /bin/bash
# 创建用户
rabbitmqctl add_user admin admin123
rabbitmqctl add_user developer dev123
rabbitmqctl add_user guest guest
# 设置用户角色
rabbitmqctl set_user_tags admin administrator
rabbitmqctl set_user_tags developer management
rabbitmqctl set_user_tags guest monitoring
# 设置权限(vhost, 配置权限, 写权限, 读权限)
rabbitmqctl set_permissions -p / admin ".*" ".*" ".*"
rabbitmqctl set_permissions -p / developer ".*" ".*" ".*"
rabbitmqctl set_permissions -p / guest ".*" ".*" ".*"
# 查看用户列表
rabbitmqctl list_users
# 查看用户权限
rabbitmqctl list_permissions -p /
虚拟主机管理
# 创建虚拟主机
rabbitmqctl add_vhost /dev
rabbitmqctl add_vhost /prod
# 列出虚拟主机
rabbitmqctl list_vhosts
# 删除虚拟主机
rabbitmqctl delete_vhost /test
# 为用户分配虚拟主机权限
rabbitmqctl set_permissions -p /dev developer ".*" ".*" ".*"
2. 队列管理
# 声明队列
rabbitmqctl declare queue name=task_queue durable=true
# 查看队列信息
rabbitmqctl list_queues name messages consumers
# 查看队列详细信息
rabbitmqctl list_queues name messages_ready messages_unacknowledged
# 删除队列
rabbitmqctl delete_queue task_queue
3. 交换机管理
# 创建交换机
rabbitmqctl declare exchange name=logs type=fanout
# 查看交换机
rabbitmqctl list_exchanges name type
# 绑定队列到交换机
rabbitmqctl bind_queue task_queue logs
# 解绑
rabbitmqctl unbind_queue task_queue logs
消息模式详解
1. 简单队列 (Simple Queue)
# 生产者 (producer.py)
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明队列
channel.queue_declare(queue='hello')
# 发送消息
channel.basic_publish(exchange='', routing_key='hello', body='Hello World!')
print(" [x] Sent 'Hello World!'")
connection.close()
# 消费者 (consumer.py)
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明队列
channel.queue_declare(queue='hello')
def callback(ch, method, properties, body):
print(f" [x] Received {body.decode()}")
# 消费消息
channel.basic_consume(queue='hello', on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
2. 工作队列 (Work Queue)
# 生产者 (new_task.py)
import pika
import sys
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明持久化队列
channel.queue_declare(queue='task_queue', durable=True)
message = ' '.join(sys.argv[1:]) or "Hello World!"
# 发送持久化消息
channel.basic_publish(
exchange='',
routing_key='task_queue',
body=message,
properties=pika.BasicProperties(
delivery_mode=2, # 消息持久化
)
)
print(f" [x] Sent {message}")
connection.close()
# 消费者 (worker.py)
import pika
import time
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
channel.queue_declare(queue='task_queue', durable=True)
def callback(ch, method, properties, body):
print(f" [x] Received {body.decode()}")
time.sleep(body.count(b'.')) # 模拟处理时间
print(" [x] Done")
ch.basic_ack(delivery_tag=method.delivery_tag) # 手动确认
# 设置公平分发
channel.basic_qos(prefetch_count=1)
channel.basic_consume(queue='task_queue', on_message_callback=callback)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
3. 发布/订阅 (Publish/Subscribe)
# 发布者 (emit_log.py)
import pika
import sys
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明 fanout 交换机
channel.exchange_declare(exchange='logs', exchange_type='fanout')
message = ' '.join(sys.argv[1:]) or "Hello World!"
# 发布消息
channel.basic_publish(exchange='logs', routing_key='', body=message)
print(f" [x] Sent {message}")
connection.close()
# 订阅者 (receive_logs.py)
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
channel.exchange_declare(exchange='logs', exchange_type='fanout')
# 创建临时队列
result = channel.queue_declare(queue='', exclusive=True)
queue_name = result.method.queue
# 绑定队列到交换机
channel.queue_bind(exchange='logs', queue=queue_name)
def callback(ch, method, properties, body):
print(f" [x] {body.decode()}")
channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for logs. To exit press CTRL+C')
channel.start_consuming()
4. 路由模式 (Routing)
# 发送者 (emit_log_direct.py)
import pika
import sys
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明 direct 交换机
channel.exchange_declare(exchange='direct_logs', exchange_type='direct')
severity = sys.argv[1] if len(sys.argv) > 1 else 'info'
message = ' '.join(sys.argv[2:]) or 'Hello World!'
# 使用路由键发送消息
channel.basic_publish(
exchange='direct_logs',
routing_key=severity,
body=message
)
print(f" [x] Sent [{severity}] {message}")
connection.close()
编程接口
1. Python 客户端
安装依赖
pip install pika
完整示例
# rabbitmq_client.py
import pika
import json
import uuid
from datetime import datetime
class RabbitMQClient:
def __init__(self, host='localhost', port=5672, username='admin', password='admin123'):
self.connection_params = pika.ConnectionParameters(
host=host,
port=port,
credentials=pika.PlainCredentials(username, password)
)
self.connection = None
self.channel = None
def connect(self):
"""建立连接"""
try:
self.connection = pika.BlockingConnection(self.connection_params)
self.channel = self.connection.channel()
print("Connected to RabbitMQ successfully")
except Exception as e:
print(f"Failed to connect to RabbitMQ: {e}")
def declare_queue(self, queue_name, durable=True):
"""声明队列"""
self.channel.queue_declare(queue=queue_name, durable=durable)
def declare_exchange(self, exchange_name, exchange_type='direct'):
"""声明交换机"""
self.channel.exchange_declare(
exchange=exchange_name,
exchange_type=exchange_type
)
def publish_message(self, exchange, routing_key, message, persistent=True):
"""发布消息"""
properties = pika.BasicProperties(
delivery_mode=2 if persistent else 1,
message_id=str(uuid.uuid4()),
timestamp=int(datetime.now().timestamp())
)
self.channel.basic_publish(
exchange=exchange,
routing_key=routing_key,
body=json.dumps(message) if isinstance(message, dict) else message,
properties=properties
)
def consume_messages(self, queue_name, callback, auto_ack=False):
"""消费消息"""
self.channel.basic_qos(prefetch_count=1)
self.channel.basic_consume(
queue=queue_name,
on_message_callback=callback,
auto_ack=auto_ack
)
print(' [*] Waiting for messages. To exit press CTRL+C')
try:
self.channel.start_consuming()
except KeyboardInterrupt:
print(' [!] Interrupted')
self.channel.stop_consuming()
def close(self):
"""关闭连接"""
if self.connection and not self.connection.is_closed:
self.connection.close()
# 使用示例
def message_handler(ch, method, properties, body):
"""消息处理函数"""
try:
message = json.loads(body)
print(f"Processing message: {message}")
# 处理消息逻辑
# ...
# 手动确认
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
print(f"Error processing message: {e}")
# 拒绝消息并重新队列
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
if __name__ == "__main__":
client = RabbitMQClient()
client.connect()
# 声明队列和交换机
client.declare_exchange('user_events', 'topic')
client.declare_queue('user_registration')
# 绑定队列到交换机
client.channel.queue_bind(
exchange='user_events',
queue='user_registration',
routing_key='user.registration'
)
# 发布消息
message = {
'user_id': 123,
'email': 'user@example.com',
'timestamp': datetime.now().isoformat()
}
client.publish_message('user_events', 'user.registration', message)
# 消费消息
client.consume_messages('user_registration', message_handler)
client.close()
2. Java 客户端
Maven 依赖
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.20.0</version>
</dependency>
完整示例
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
public class RabbitMQClient {
private Connection connection;
private Channel channel;
public void connect() throws IOException, TimeoutException {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("admin123");
connection = factory.newConnection();
channel = connection.createChannel();
}
public void publishMessage(String exchange, String routingKey, String message)
throws IOException {
channel.basicPublish(exchange, routingKey,
MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes());
}
public void consumeMessages(String queueName) throws IOException {
channel.basicQos(1);
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
System.out.println(" [x] Received '" + message + "'");
try {
// 处理消息
processMessage(message);
// 确认消息
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
} catch (Exception e) {
System.err.println("Error processing message: " + e.getMessage());
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
}
};
channel.basicConsume(queueName, false, deliverCallback, consumerTag -> {});
}
private void processMessage(String message) {
// 消息处理逻辑
System.out.println("Processing: " + message);
}
public void close() throws IOException, TimeoutException {
if (channel != null && channel.isOpen()) {
channel.close();
}
if (connection != null && connection.isOpen()) {
connection.close();
}
}
}
高级特性
1. 消息确认机制
# 生产者确认
import pika
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 启用发布确认
channel.confirm_delivery()
try:
channel.basic_publish(
exchange='',
routing_key='test_queue',
body='Hello World!',
mandatory=True
)
print("Message delivered successfully")
except pika.exceptions.UnroutableError:
print("Message was returned")
2. 死信队列
# 设置死信队列
channel.queue_declare(
queue='main_queue',
durable=True,
arguments={
'x-message-ttl': 60000, # 消息TTL
'x-dead-letter-exchange': 'dlx_exchange',
'x-dead-letter-routing-key': 'dlx_routing_key'
}
)
# 声明死信交换机和队列
channel.exchange_declare(exchange='dlx_exchange', exchange_type='direct')
channel.queue_declare(queue='dead_letter_queue', durable=True)
channel.queue_bind(exchange='dlx_exchange', queue='dead_letter_queue', routing_key='dlx_routing_key')
3. 优先级队列
# 声明优先级队列
channel.queue_declare(
queue='priority_queue',
durable=True,
arguments={'x-max-priority': 10}
)
# 发送优先级消息
channel.basic_publish(
exchange='',
routing_key='priority_queue',
body='High priority message',
properties=pika.BasicProperties(priority=9)
)
监控与管理
1. Web 管理界面
访问管理界面:
- URL: http://localhost:15672
- 用户名: admin
- 密码: admin123
管理界面功能:
- Overview: 集群状态概览
- Connections: 连接管理
- Channels: 通道管理
- Exchanges: 交换机管理
- Queues: 队列管理
- Admin: 用户和权限管理
2. 命令行监控
# 系统状态
rabbitmqctl status
# 集群状态
rabbitmqctl cluster_status
# 内存使用情况
rabbitmqctl environment
# 连接信息
rabbitmqctl list_connections
# 通道信息
rabbitmqctl list_channels
# 队列信息
rabbitmqctl list_queues name messages consumers memory
# 交换机信息
rabbitmqctl list_exchanges name type
# 绑定信息
rabbitmqctl list_bindings
3. 性能监控
# 启用性能监控插件
rabbitmq-plugins enable rabbitmq_management
# 查看节点健康状态
rabbitmq-diagnostics ping
rabbitmq-diagnostics check_running
rabbitmq-diagnostics check_local_alarms
# 监控队列性能
rabbitmqctl list_queues name messages_ready messages_unacknowledged consumers message_stats.publish_details.rate
性能优化
1. 连接池配置
# 连接池实现
import pika
from queue import Queue
import threading
class RabbitMQConnectionPool:
def __init__(self, host='localhost', port=5672, username='admin', password='admin123', pool_size=10):
self.connection_params = pika.ConnectionParameters(
host=host,
port=port,
credentials=pika.PlainCredentials(username, password)
)
self.pool = Queue(maxsize=pool_size)
self.pool_size = pool_size
self.lock = threading.Lock()
# 初始化连接池
for _ in range(pool_size):
connection = pika.BlockingConnection(self.connection_params)
self.pool.put(connection)
def get_connection(self):
"""获取连接"""
return self.pool.get()
def return_connection(self, connection):
"""归还连接"""
if not connection.is_closed:
self.pool.put(connection)
2. 批量处理
# 批量发送消息
def batch_publish(channel, exchange, routing_key, messages, batch_size=100):
"""批量发送消息"""
for i in range(0, len(messages), batch_size):
batch = messages[i:i+batch_size]
for message in batch:
channel.basic_publish(
exchange=exchange,
routing_key=routing_key,
body=message
)
# 等待确认
channel.confirm_delivery()
3. 预取设置
# 设置预取数量
channel.basic_qos(prefetch_count=100) # 预取100条消息
# 设置全局预取
channel.basic_qos(prefetch_count=100, global_qos=True)
常见问题
问题1:管理界面无法访问
现象: 访问 http://localhost:15672 返回 404 原因: 管理插件未启用或容器端口映射错误 解决方案:
# 启用管理插件
docker exec -it rabbitmq rabbitmq-plugins enable rabbitmq_management
# 检查端口映射
docker port rabbitmq
问题2:消息丢失
现象: 发送的消息没有被消费者接收 原因: 队列或消息未持久化 解决方案:
# 声明持久化队列
channel.queue_declare(queue='my_queue', durable=True)
# 发送持久化消息
channel.basic_publish(
exchange='',
routing_key='my_queue',
body='Hello World!',
properties=pika.BasicProperties(delivery_mode=2)
)
问题3:消息积压
现象: 队列中消息数量持续增长 原因: 消费者处理速度跟不上生产速度 解决方案:
# 增加消费者数量
# 调整预取数量
channel.basic_qos(prefetch_count=1)
# 使用工作队列模式
# 优化消息处理逻辑
问题4:内存不足
现象: RabbitMQ 进程内存使用过高 原因: 消息积压或内存配置不当 解决方案:
# 设置内存限制
docker run -d --memory=2g --name rabbitmq rabbitmq:3.12-management
# 配置内存告警
echo "vm_memory_high_watermark.relative = 0.6" >> /etc/rabbitmq/rabbitmq.conf
问题5:集群脑裂
现象: 集群节点间通信异常 原因: 网络分区或节点故障 解决方案:
# 检查集群状态
rabbitmqctl cluster_status
# 处理网络分区
rabbitmqctl force_boot
# 重新加入集群
rabbitmqctl join_cluster rabbit@node1
⚠️ 注意: 在生产环境中,建议使用集群模式部署以确保高可用性,并定期监控队列积压情况。
问题6:管理界面不显示统计数据
现象: 管理界面能打开,但 Overview、Queues 页面的消息速率和图表为空,提示 “Stats in management UI are disabled on this node”
原因: 官方 Docker 镜像默认在 /etc/rabbitmq/conf.d/ 中关闭了指标收集(management_agent.disable_metrics_collector = true)
解决方案:
# 进入容器,覆盖该配置
docker exec -it rabbitmq bash
cd /etc/rabbitmq/conf.d/
echo management_agent.disable_metrics_collector = false > management_agent.disable_metrics_collector.conf
exit
# 重启容器生效
docker restart rabbitmq
相关文章
消息队列对比
相关中间件
数据库技术
大数据技术
运维部署
监控运维
相关阅读
- Kafka 简明教程:
./Kafka简明教程.md - Redis 入门:
./Redis.md - Nacos 配置中心:
./Nacos.md - Kubernetes 本地部署:
../../../运维/Linux/容器编排/Kubernetes/本地部署K8s集群(一键部署).md