VB.net 2010 视频教程 VB.net 2010 视频教程 python基础视频教程
SQL Server 2008 视频教程 c#入门经典教程 Visual Basic从门到精通视频教程
当前位置:
首页 > 编程开发 > c#编程 >
  • 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的不足,适合需要可靠消息的实时场景(比如聊天、通知)。
适用场景:简单的实时聊天、通知、实时监控等不需要持久化的场景。
六、三种消息队列对比与最佳实践

  1. 三种消息队列对比
    特性 RabbitMQ Kafka Redis Pub/Sub
    可靠性 高(持久化、手动确认) 高(持久化、副本) 低(不持久化)
    吞吐量 中(万级QPS) 高(百万级QPS) 中(十万级QPS)
    功能丰富度 高(交换机、死信队列、延迟队列) 中(分区、副本、流处理) 低(仅发布订阅)
    适用场景 复杂业务场景(订单、支付、物流) 大数据场景(日志、实时计算) 简单实时场景(聊天、通知)
    部署复杂度 中(单节点或集群) 高(分布式集群) 低(单节点或集群)
    .NET生态支持 好(RabbitMQ.Client成熟) 好(Confluent.Kafka成熟) 好(StackExchange.Redis成熟)
  2. 最佳实践
    (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

相关教程