全部笔记All notes

RabbitMQ 消息队列完整指南

阅读 9m 17s9m 17s read

概述

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 管理界面

访问管理界面:

管理界面功能:

  • 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