RabbitMQ 消息队列
2024/4/15大约 4 分钟
RabbitMQ 消息队列
概述
RabbitMQ 是一个开源的消息代理软件,实现了 AMQP(高级消息队列协议),用于在分布式系统中进行异步通信。
1. 安装与配置
1.1 使用 Docker 安装
# 拉取 RabbitMQ 镜像
docker pull rabbitmq:3-management
# 启动容器
docker run -d \
--name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=password \
rabbitmq:3-management1.2 访问管理界面
打开浏览器访问 http://localhost:15672,使用用户名 admin 和密码 password 登录。
2. 核心概念
2.1 消息模型
生产者 -> 交换机 -> 队列 -> 消费者2.2 交换机类型
| 类型 | 说明 |
|---|---|
| Direct | 精确匹配路由键 |
| Fanout | 广播到所有绑定的队列 |
| Topic | 模式匹配路由键 |
| Headers | 根据消息头匹配 |
3. .NET 客户端使用
3.1 安装依赖
dotnet add package RabbitMQ.Client3.2 生产者实现
using RabbitMQ.Client;
using System.Text;
var factory = new ConnectionFactory()
{
HostName = "localhost",
UserName = "admin",
Password = "password"
};
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
// 声明队列
channel.QueueDeclare(queue: "hello",
durable: false,
exclusive: false,
autoDelete: false,
arguments: null);
string message = "Hello RabbitMQ!";
var body = Encoding.UTF8.GetBytes(message);
// 发布消息
channel.BasicPublish(exchange: "",
routingKey: "hello",
basicProperties: null,
body: body);
Console.WriteLine($"发送消息: {message}");
}
Console.WriteLine("按任意键退出...");
Console.ReadKey();3.3 消费者实现
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System.Text;
var factory = new ConnectionFactory()
{
HostName = "localhost",
UserName = "admin",
Password = "password"
};
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
// 声明队列(与生产者一致)
channel.QueueDeclare(queue: "hello",
durable: false,
exclusive: false,
autoDelete: false,
arguments: null);
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
Console.WriteLine($"收到消息: {message}");
};
// 消费消息
channel.BasicConsume(queue: "hello",
autoAck: true,
consumer: consumer);
Console.WriteLine("等待消息...");
Console.ReadKey();
}4. 高级特性
4.1 消息持久化
// 声明持久化队列
channel.QueueDeclare(queue: "task_queue",
durable: true, // 队列持久化
exclusive: false,
autoDelete: false,
arguments: null);
// 发送持久化消息
var properties = channel.CreateBasicProperties();
properties.Persistent = true; // 消息持久化
channel.BasicPublish(exchange: "",
routingKey: "task_queue",
basicProperties: properties,
body: body);4.2 消息确认机制
// 关闭自动确认
channel.BasicConsume(queue: "task_queue",
autoAck: false, // 手动确认
consumer: consumer);
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
Console.WriteLine($"收到消息: {message}");
// 模拟处理时间
Thread.Sleep(1000);
Console.WriteLine("消息处理完成");
// 手动确认消息
channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
};4.3 公平分发
// 设置每次只接收一条未确认的消息
channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);5. 交换机使用
5.1 Direct 交换机
// 声明 Direct 交换机
channel.ExchangeDeclare(exchange: "direct_logs", type: ExchangeType.Direct);
// 绑定队列到交换机
channel.QueueBind(queue: "queue1", exchange: "direct_logs", routingKey: "error");
channel.QueueBind(queue: "queue2", exchange: "direct_logs", routingKey: "info");
channel.QueueBind(queue: "queue2", exchange: "direct_logs", routingKey: "warning");
// 发送消息到交换机
channel.BasicPublish(exchange: "direct_logs",
routingKey: "error",
basicProperties: null,
body: body);5.2 Fanout 交换机
// 声明 Fanout 交换机
channel.ExchangeDeclare(exchange: "fanout_logs", type: ExchangeType.Fanout);
// 绑定队列到交换机(忽略路由键)
channel.QueueBind(queue: "queue1", exchange: "fanout_logs", routingKey: "");
channel.QueueBind(queue: "queue2", exchange: "fanout_logs", routingKey: "");
// 发送消息(路由键会被忽略)
channel.BasicPublish(exchange: "fanout_logs",
routingKey: "",
basicProperties: null,
body: body);5.3 Topic 交换机
// 声明 Topic 交换机
channel.ExchangeDeclare(exchange: "topic_logs", type: ExchangeType.Topic);
// 绑定队列(使用通配符)
channel.QueueBind(queue: "queue1", exchange: "topic_logs", routingKey: "*.orange.*");
channel.QueueBind(queue: "queue2", exchange: "topic_logs", routingKey: "*.*.rabbit");
channel.QueueBind(queue: "queue2", exchange: "topic_logs", routingKey: "lazy.#");
// 发送消息
channel.BasicPublish(exchange: "topic_logs",
routingKey: "quick.orange.rabbit",
basicProperties: null,
body: body);6. 死信队列
6.1 创建死信交换机和队列
// 声明死信交换机
channel.ExchangeDeclare(exchange: "dlx_exchange", type: ExchangeType.Direct);
// 声明死信队列
channel.QueueDeclare(queue: "dlx_queue", durable: true, exclusive: false, autoDelete: false);
// 绑定死信队列到死信交换机
channel.QueueBind(queue: "dlx_queue", exchange: "dlx_exchange", routingKey: "#");
// 声明普通队列,并设置死信交换机
var args = new Dictionary<string, object>
{
{ "x-dead-letter-exchange", "dlx_exchange" },
{ "x-dead-letter-routing-key", "#" }
};
channel.QueueDeclare(queue: "normal_queue", durable: true, exclusive: false, autoDelete: false, arguments: args);6.2 消息过期
// 设置消息过期时间(毫秒)
var properties = channel.CreateBasicProperties();
properties.Expiration = "60000"; // 60秒过期
channel.BasicPublish(exchange: "",
routingKey: "normal_queue",
basicProperties: properties,
body: body);7. 实践建议
7.1 最佳实践
- 队列命名规范:使用有意义的队列名称
- 消息持久化:生产环境建议开启持久化
- 消息确认:使用手动确认保证消息不丢失
- 消息过期:设置合理的消息过期时间
- 监控告警:监控队列长度和消息堆积
7.2 常见问题
| 问题 | 解决方案 |
|---|---|
| 消息丢失 | 开启持久化、使用手动确认 |
| 消息重复 | 消息去重、使用幂等性设计 |
| 队列堆积 | 增加消费者、优化消费速度 |
| 网络问题 | 设置合理的超时时间和重试机制 |
总结
RabbitMQ 是一个功能强大的消息队列系统,支持多种消息模式和高级特性。合理使用可以提升系统的异步处理能力和可靠性。
作者:Blogger
日期:2024年4月15日