-
C#网络编程之消息队列(RabbitMQ、Kafka、Redis Pub/Sub)
第28章 消息队列实战
28.1 消息队列(RabbitMQ、Kafka、Redis Pub/Sub)
一、我踩过的消息队列坑:从“同步调用雪崩”到“削峰填谷救系统”
刚做电商系统时,订单创建后同步调用库存服务、支付服务、物流服务,结果库存服务挂了,整个订单系统直接雪崩,1小时损失5000单;后来用RabbitMQ做异步调用,订单创建后发一条消息到队列,库存服务异步消费,即使库存服vb.net教程C#教程python教程SQL教程access 2010教程务挂了,消息也存在队列里,恢复后自动处理,系统稳定性直接提升到99.99%;但又踩了一堆坑:用Redis Pub/Sub做聊天功能,结果订阅者离线后错过消息,用户投诉聊天记录丢失;用Kafka做日志收集,分区数设成1,导致单线程处理,CPU跑满,日志堆积了100G;用RabbitMQ时忘记确认消息,导致消费者重复消费,库存扣了两次,用户收到两个包裹。这节我把这些踩坑经验揉进去,用大白话讲透三种主流消息队列的核心原理,结合C#实战代码逐行讲解,对比它们的优缺点和适用场景,以及消息幂等性、重试机制等最佳实践,让你的异步通信既可靠又高效。
二、消息队列核心作用:解耦、异步、削峰填谷(拓展知识)
消息队列是一种异步通信模式,核心作用有三个:
1.解耦:系统之间通过消息队列通信,不需要直接调用,比如订单系统不需要知道库存服务的地址,只需要发消息到队列;
2.异步:不需要等待服务处理完成,比如订单创建后直接返回成功,库存服务异步处理,提升用户体验;
3.削峰填谷:把突发的请求(比如秒杀活动的10万订单)存入队列,系统慢慢处理,避免服务器崩溃。
类比:消息队列就像你去快递站寄快递,不需要等快递员送到收件人手里,只需要把快递放在快递站,快递员会异步处理,快递站还能存很多快递,避免你排队等待。
三、RabbitMQ:可靠的消息中间件,适合复杂业务场景
RabbitMQ是基于AMQP协议的消息中间件,核心优势是可靠性高、功能丰富,支持交换机、队列、绑定、死信队列、延迟队列等功能,适合复杂的业务场景(比如订单、支付、物流)。
核心概念(大白话)
1.生产者:发消息的服务(比如订单系统);
2.消费者:收消息的服务(比如库存系统);
3.交换机(Exchange):接收生产者的消息,根据路由规则转发到队列;
4.队列(Queue):存储消息,消费者从队列里取消息;
5.绑定(Binding):把交换机和队列绑定,指定路由规则。
类比:生产者是寄快递的人,交换机是快递站的分拣员,队列是不同地区的快递筐,绑定是分拣规则(比如北京的快递放到北京筐),消费者是快递员,从筐里取快递送出去。
实战1:RabbitMQ生产者与消费者(C#实战)
用RabbitMQ.Client库实现订单系统发消息,库存系统消费消息。
步骤1:安装RabbitMQ.Client NuGet包
bash
Install-Package RabbitMQ.Client
步骤2:生产者代码(订单系统)
csharp
using RabbitMQ.Client;
using System;
using System.Text;
namespace RabbitMQProducer;
class Program
{
static void Main(string[] args)
{
// 1. 创建连接工厂
var factory = new ConnectionFactory
{
HostName = "localhost", // RabbitMQ地址
UserName = "guest", // 默认用户名
Password = "guest" // 默认密码
};
// 2. 创建连接和通道
using var connection = factory.CreateConnection();
using var channel = connection.CreateModel();
// 3. 声明交换机(持久化,避免重启后丢失)
string exchangeName = "order_exchange";
channel.ExchangeDeclare(exchange: exchangeName, type: ExchangeType.Direct, durable: true);
// 4. 声明队列(持久化,避免重启后丢失)
string queueName = "stock_queue";
channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
// 5. 绑定交换机和队列,指定路由键(比如"order.created")
string routingKey = "order.created";
channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: routingKey);
// 6. 构造消息(订单ID)
string orderId = Guid.NewGuid().ToString();
var body = Encoding.UTF8.GetBytes(orderId);
// 7. 发消息(持久化,避免RabbitMQ重启后丢失)
channel.BasicPublish(
exchange: exchangeName,
routingKey: routingKey,
basicProperties: channel.CreateBasicProperties { Persistent = true }, // 持久化消息
body: body);
Console.WriteLine($"[生产者] 已发送订单消息:{orderId}");
Console.WriteLine("Press any key to exit...");
Console.ReadKey();
}
}
代码逐行讲解:
1.ConnectionFactory:配置RabbitMQ的连接信息,生产环境要从配置文件读取;
2.ExchangeDeclare:声明交换机,durable: true表示持久化,RabbitMQ重启后交换机不会丢失;
3.QueueDeclare:声明队列,durable: true表示持久化队列,exclusive: false表示队列可以被多个消费者访问;
4.QueueBind:绑定交换机和队列,routingKey是路由规则,只有匹配的消息才会被转发到队列;
5.BasicPublish:发消息,Persistent: true表示持久化消息,RabbitMQ重启后消息不会丢失;
6.using语句:自动释放连接和通道,避免资源泄漏。
我踩过的坑:一开始没设置Persistent: true,RabbitMQ重启后消息全部丢失,导致订单和库存不一致——生产环境必须设置消息持久化!
步骤3:消费者代码(库存系统)
csharp
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System;
using System.Text;
using System.Threading;
namespace RabbitMQConsumer;
class Program
{
static void Main(string[] args)
{
var factory = new ConnectionFactory
{
HostName = "localhost",
UserName = "guest",
Password = "guest"
};
using var connection = factory.CreateConnection();
using var channel = connection.CreateModel();
string exchangeName = "order_exchange";
string queueName = "stock_queue";
string routingKey = "order.created";
// 声明交换机、队列、绑定(和生产者一致,避免消费者先启动时不存在)
channel.ExchangeDeclare(exchange: exchangeName, type: ExchangeType.Direct, durable: true);
channel.QueueDeclare(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null);
channel.QueueBind(queue: queueName, exchange: exchangeName, routingKey: routingKey);
// 配置消费者:每次只取1条消息,处理完成后再取(避免消费者崩溃导致消息丢失)
channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);
// 创建消费者
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var orderId = Encoding.UTF8.GetString(body);
Console.WriteLine($"[消费者] 收到订单消息:{orderId}");
// 模拟库存扣减(耗时操作)
Thread.Sleep(1000);
Console.WriteLine($"[消费者] 已完成订单{orderId}的库存扣减");
// 确认消息处理完成(必须!否则RabbitMQ会重新发送消息)
channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
};
// 启动消费者(自动确认:false,手动确认消息)
channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
Console.WriteLine("[消费者] 等待订单消息...");
Console.WriteLine("Press any key to exit...");
Console.ReadKey();
}
}
代码逐行讲解:
1.BasicQos:设置每次只取1条消息,处理完成后再取,避免消费者崩溃导致大量消息未处理;
2.EventingBasicConsumer:事件驱动的消费者,Received事件触发时处理消息;
3.BasicAck:手动确认消息处理完成,RabbitMQ会把消息从队列里删除;
1.如果忘记确认,RabbitMQ会认为消息未处理,消费者重启后会重新发送消息,导致重复消费;
4.autoAck: false:关闭自动确认,必须手动确认消息,保证消息不丢失。
我踩过的坑:一开始用autoAck: true,消费者处理消息时崩溃,消息被RabbitMQ删除,导致库存没扣减,订单和库存不一致——生产环境必须用手动确认!
拓展知识:RabbitMQ常用功能
1.死信队列:处理消费失败的消息,比如库存扣减失败3次后,把消息放到死信队列,人工处理;
2.延迟队列:实现延迟任务,比如订单创建后30分钟未支付,自动取消订单;
3.交换机类型:
1.Direct:精确匹配路由键,适合一对一通信;
2.Fanout:把消息转发到所有绑定的队列,适合广播(比如通知所有服务更新配置);
3.Topic:模糊匹配路由键(比如order.*匹配order.created、order.paid),适合一对多通信。
适用场景:订单系统、支付系统、物流系统等需要可靠消息的复杂业务场景。
四、Kafka:高性能的消息中间件,适合大数据场景
Kafka是基于分布式发布订阅的消息系统,核心优势是高性能、高吞吐量,支持百万级QPS,适合大数据场景(比如日志收集、实时计算、用户行为分析)。
核心概念(大白话)
1.主题(Topic):消息的分类,比如“order_topic”、“log_topic”;
2.分区(Partition):每个主题分成多个分区,消息按顺序存储在分区里,分区越多,吞吐量越高;
3.副本(Replica):每个分区有多个副本,保证高可用,比如一个分区有3个副本,其中1个是主副本,2个是从副本;
4.生产者:发消息到主题的分区;
5.消费者组:多个消费者组成一个组,每个分区只能被组内的一个消费者消费,保证消息不重复消费。
类比:主题是一个大仓库,分区是仓库里的货架,副本是货架的备份,生产者是把货物放到货架上的人,消费者组是一群工人,每个货架只能被一个工人取货,保证货物不被重复取走。
实战2:Kafka生产者与消费者(C#实战)
用Confluent.Kafka库实现日志收集系统,生产者发日志到Kafka,消费者消费日志并存入Elasticsearch。
步骤1:安装Confluent.Kafka NuGet包
bash
Install-Package Confluent.Kafka
步骤2:生产者代码(日志系统)
csharp
using Confluent.Kafka;
using System;
namespace KafkaProducer;
class Program
{
static void Main(string[] args)
{
// 1. 配置生产者
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092", // Kafka地址
Acks = Acks.All, // 等待所有副本确认消息,保证可靠性
Retries = 3, // 发送失败时重试3次
BatchSize = 16384, // 批量发送的大小(16KB)
LingerMs = 5 // 等待5ms,攒够批量再发送,提升吞吐量
};
// 2. 创建生产者
using var producer = new ProducerBuilder<Null, string>(config).Build();
try
{
// 3. 发消息到"log_topic"主题
var topic = "log_topic";
var message = $"[{DateTime.Now:yyyy-MM-dd HH:mm:ss}] INFO: User 123 logged in";
var deliveryResult = producer.ProduceAsync(topic, new Message<Null, string> { Value = message }).Result;
Console.WriteLine($"[生产者] 已发送日志到分区{deliveryResult.Partition},偏移量{deliveryResult.Offset}:{message}");
}
catch (ProduceException<Null, string> e)
{
Console.WriteLine($"[生产者] 发送失败:{e.Error.Reason}");
}
Console.WriteLine("Press any key to exit...");
Console.ReadKey();
}
}
代码逐行讲解:
1.ProducerConfig:配置Kafka生产者,Acks.All表示等待所有副本确认消息,保证消息不丢失;
2.ProducerBuilder:创建生产者,Null表示键为空,string表示值是字符串;
3.ProduceAsync:异步发消息到主题,返回DeliveryResult,包含分区和偏移量;
4.BatchSize和LingerMs:批量发送配置,攒够16KB或等待5ms再发送,提升吞吐量。
我踩过的坑:一开始把Acks设为Acks.None,Kafka服务器崩溃后消息全部丢失,生产环境必须设为Acks.All!
步骤3:消费者代码(日志收集系统)
csharp
using Confluent.Kafka;
using System;
namespace KafkaConsumer;
class Program
{
static void Main(string[] args)
{
// 1. 配置消费者
var config = new ConsumerConfig
{
BootstrapServers = "localhost:9092",
GroupId = "log_consumer_group", // 消费者组ID,同一组的消费者共享分区
AutoOffsetReset = AutoOffsetReset.Earliest, // 消费者启动时,从最早的消息开始消费
EnableAutoCommit = false, // 关闭自动提交偏移量,手动提交
AutoCommitIntervalMs = 5000 // 如果开启自动提交,每5秒提交一次
};
// 2. 创建消费者
using var consumer = new ConsumerBuilder<Ignore, string>(config).Build();
// 3. 订阅"log_topic"主题
var topic = "log_topic";
consumer.Subscribe(topic);
Console.WriteLine("[消费者] 等待日志消息...");
try
{
while (true)
{
// 4. 拉取消息(超时1秒)
var consumeResult = consumer.Consume(TimeSpan.FromSeconds(1));
if (consumeResult == null)
continue;
Console.WriteLine($"[消费者] 从分区{consumeResult.Partition},偏移量{consumeResult.Offset}收到日志:{consumeResult.Message.Value}");
// 模拟存入Elasticsearch
Console.WriteLine($"[消费者] 已将日志存入Elasticsearch");
// 5. 手动提交偏移量(保证消息处理完成后再提交)
consumer.Commit(consumeResult);
}
}
catch (ConsumeException e)
{
Console.WriteLine($"[消费者] 消费失败:{e.Error.Reason}");
}
finally
{
consumer.Close();
}
}
}
代码逐行讲解:
1.ConsumerConfig:配置Kafka消费者,GroupId是消费者组ID,同一组的消费者共享主题的分区;
2.AutoOffsetReset:消费者启动时,如果没有偏移量,从最早的消息开始消费(Earliest),或者从最新的消息开始消费(Latest);
3.EnableAutoCommit:关闭自动提交偏移量,手动提交,保证消息处理完成后再提交,避免消息丢失;
4.Consume:拉取消息,超时1秒,避免阻塞;
5.Commit:手动提交偏移量,Kafka会记录消费者组的偏移量,下次启动时从偏移量的下一条消息开始消费。
我踩过的坑:一开始把GroupId设成一样的,启动了3个消费者,结果每个消费者消费所有分区的消息,导致重复消费——同一消费者组的消费者,每个分区只能被一个消费者消费,不同组的消费者可以消费同一分区的消息!
拓展知识:Kafka常用功能
1.分区策略:生产者可以指定键,相同键的消息会被发到同一个分区,保证消息的顺序性;
2.流处理:用Kafka Streams实现实时流处理,比如计算实时订单量、用户活跃度;
3.连接器:用Kafka Connect把Kafka和其他系统集成,比如从MySQL同步数据到Kafka,从Kafka同步数据到Elasticsearch。
适用场景:日志收集、实时计算、用户行为分析、大数据传输等需要高吞吐量的场景。
五、Redis Pub/Sub:轻量级的消息中间件,适合实时场景
Redis Pub/Sub是Redis的发布订阅功能,核心优势是轻量级、实时性高,适合简单的实时场景(比如聊天、通知、实时监控)。
核心概念(大白话)
1.发布者:发消息到频道(Channel);
2.订阅者:订阅频道,接收频道的消息;
3.频道:消息的分类,比如“chat_channel”、“monitor_channel”。
类比:发布者是电台主持人,频道是电台频率,订阅者是听众,调对频率就能听到主持人的消息。
实战3:Redis Pub/Sub发布者与订阅者(C#实战)
用StackExchange.Redis库实现简单的聊天功能,发布者发消息到频道,订阅者接收消息。
步骤1:安装StackExchange.Redis NuGet包
bash
Install-Package StackExchange.Redis
步骤2:发布者代码(聊天客户端)
csharp
using StackExchange.Redis;
using System;
namespace RedisPublisher;
class Program
{
static void Main(string[] args)
{
// 1. 连接Redis
var redis = ConnectionMultiplexer.Connect("localhost:6379");
var db = redis.GetDatabase();
var channel = "chat_channel";
Console.WriteLine("Enter your message (or 'exit' to quit):");
while (true)
{
var message = Console.ReadLine();
if (message == "exit")
break;
// 2. 发消息到频道
db.Publish(channel, message);
Console.WriteLine($"[发布者] 已发送消息:{message}");
}
redis.Close();
}
}
代码逐行讲解:
1.ConnectionMultiplexer:连接Redis,单例模式,避免重复创建连接;
2.Publish:发消息到指定频道,返回订阅者的数量;
3.GetDatabase:获取Redis数据库实例。
步骤3:订阅者代码(聊天客户端)
csharp
using StackExchange.Redis;
using System;
namespace RedisSubscriber;
class Program
{
static void Main(string[] args)
{
// 1. 连接Redis
var redis = ConnectionMultiplexer.Connect("localhost:6379");
var subscriber = redis.GetSubscriber();
var channel = "chat_channel";
Console.WriteLine("[订阅者] 等待聊天消息...");
// 2. 订阅频道,接收消息
subscriber.Subscribe(channel, (channel, message) =>
{
Console.WriteLine($"[订阅者] 收到消息:{message}");
});
Console.WriteLine("Press any key to exit...");
Console.ReadKey();
// 3. 取消订阅
subscriber.Unsubscribe(channel);
redis.Close();
}
}
代码逐行讲解:
1.GetSubscriber:获取Redis订阅者实例;
2.Subscribe:订阅频道,回调函数接收消息;
3.Unsubscribe:取消订阅频道。
我踩过的坑:用Redis Pub/Sub做聊天功能,订阅者离线后错过消息,用户投诉聊天记录丢失——Redis Pub/Sub不持久化消息,订阅者离线后不会收到之前的消息,适合实时性高但不需要持久化的场景!
拓展知识:Redis Pub/Sub vs Redis Stream
Redis 5.0推出了Redis Stream,支持消息持久化、消费组、消息确认,弥补了Pub/Sub的不足,适合需要可靠消息的实时场景(比如聊天、通知)。
适用场景:简单的实时聊天、通知、实时监控等不需要持久化的场景。
六、三种消息队列对比与最佳实践
-
三种消息队列对比
特性 RabbitMQ Kafka Redis Pub/Sub
可靠性 高(持久化、手动确认) 高(持久化、副本) 低(不持久化)
吞吐量 中(万级QPS) 高(百万级QPS) 中(十万级QPS)
功能丰富度 高(交换机、死信队列、延迟队列) 中(分区、副本、流处理) 低(仅发布订阅)
适用场景 复杂业务场景(订单、支付、物流) 大数据场景(日志、实时计算) 简单实时场景(聊天、通知)
部署复杂度 中(单节点或集群) 高(分布式集群) 低(单节点或集群)
.NET生态支持 好(RabbitMQ.Client成熟) 好(Confluent.Kafka成熟) 好(StackExchange.Redis成熟) -
最佳实践
(1)消息幂等性
消费者处理消息时,要保证幂等性,比如用订单ID作为唯一键,处理前检查是否已经处理过;
比如库存扣减,用订单ID作为Redis的键,处理前检查键是否存在,存在则跳过,不存在则扣减库存并设置键。
(2)重试机制
消费失败时,要重试,比如用RabbitMQ的死信队列,重试3次后放到死信队列,人工处理;
不要无限重试,避免死循环,比如重试3次后放弃。
(3)监控告警
监控消息队列的指标:队列长度、消息延迟、消费速度;
比如RabbitMQ的队列长度超过1000时,发送告警通知;Kafka的分区偏移量落后超过10000时,发送告警。
(4)选择合适的消息队列
复杂业务场景用RabbitMQ;
大数据场景用Kafka;
简单实时场景用Redis Pub/Sub或Redis Stream。
(5)生产环境部署
RabbitMQ:集群部署3-5节点,保证高可用;
Kafka:分布式集群部署,每个主题设置3个副本,分区数根据服务器CPU核心数设置(比如16核CPU设置16个分区);
Redis:集群部署,保证高可用。
七、总结
消息队列是异步通信的核心组件,解耦了系统,提升了系统的稳定性和性能。RabbitMQ适合复杂业务场景,Kafka适合大数据场景,Redis Pub/Sub适合简单实时场景。选择合适的消息队列,结合幂等性、重试机制、监控告警,能让你的系统既可靠又高效。
下一节我们会学习分布式事务:用消息队列实现最终一致性,解决跨系统的事务问题。
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49556.html










