VB.net 2010 视频教程 VB.net 2010 视频教程 python基础视频教程
SQL Server 2008 视频教程 c#入门经典教程 Visual Basic从门到精通视频教程
当前位置:
首页 > 编程开发 > c#编程 >
  • C#网络编程之消息队列(RabbitMQ、Kafka、Redis Pub/Sub)

第48章 消息队列实战
48.1 消息队列(RabbitMQ、Kafka、Redis Pub/Sub)
一、我踩过的消息队列坑:从“Redis丢单”到“Kafka数据倾斜”
做电商订单通知时,我一开始图省事用Redis Pub/Sub,结果大促时客户端断开重连,丢失了1000条订单通知;换成RabbitMQ后,又因为没开消息确认,导致消息丢失500条;后来用Kafka处理日志,结果分区键选了固定值,导致单分区过载,QPS从10000跌到2000。这节我把这些血泪经验揉进去,用大白话讲vb.net教程C#教程python教程SQL教程access 2010教程透三种主流消息队列的核心原理,结合C#实战代码逐行拆解,拓展生产级优化技巧,让你在可靠性、性能、易用性之间找到平衡。
二、Redis Pub/Sub:简单粗暴的实时通知,适合低可靠性场景
Redis Pub/Sub是Redis的发布订阅功能,核心是“广播消息”,实现简单,性能好,但不支持持久化,客户端断开会丢失消息,适合实时通知、低可靠性场景(如聊天消息、实时监控)。
核心原理(大白话+广播喇叭例子)
把Redis Pub/Sub比作“小区广播喇叭”:
1.发布者:物业用喇叭喊“停水通知”(发布消息到频道);
2.订阅者:在家的业主听到通知(订阅频道的客户端收到消息);
3.缺点:不在家的业主听不到(客户端断开后,期间的消息会丢失)。
我踩过的坑:大促时客户端断开重连,丢失了1000条订单通知——Redis Pub/Sub不支持持久化,不能用于需要可靠消息的场景!
实战1:Redis Pub/Sub实现实时通知(C#)
步骤1:安装NuGet包
bash
Install-Package StackExchange.Redis
步骤2:发布者代码(逐行讲解)
csharp

	using StackExchange.Redis;
	using System;
	using System.Threading.Tasks;
	
	namespace MessageQueue.RedisPubSub;
	
	public class RedisPublisher
	{
	private readonly IConnectionMultiplexer _redis;
	private readonly ISubscriber _subscriber;
	
	public RedisPublisher(string redisConnectionString)
	{
	// 1. 创建Redis连接复用器(生产环境用单例)
	_redis = ConnectionMultiplexer.Connect(redisConnectionString);
	// 2. 获取订阅者实例
	_subscriber = _redis.GetSubscriber();
	Console.WriteLine("Redis发布者已连接");
	}
	
	/// <summary>
	/// 发布消息到频道
	/// </summary>
	public async Task PublishAsync(string channel, string message)
	{
	// 3. 发布消息到指定频道
	await _subscriber.PublishAsync(channel, message);
	Console.WriteLine($"发布消息到频道 {channel}:{message}");
	}
	
	public void Dispose()
	{
	_redis.Dispose();
	}
	}
	
	// 测试代码:模拟订单服务发布通知
	class Program
	{
	static async Task Main(string[] args)
	{
	var publisher = new RedisPublisher("localhost:6379");
	for (int i = 0; i < 10; i++)
	{
	string message = $"订单 {i+1} 已支付";
	await publisher.PublishAsync("order_notify", message);
	await Task.Delay(1000); // 每秒发一条
	}
	publisher.Dispose();
	}
	}

步骤3:订阅者代码(逐行讲解)
csharp

	using StackExchange.Redis;
	using System;
	using System.Threading.Tasks;
	
	namespace MessageQueue.RedisPubSub;
	
	public class RedisSubscriber
	{
	private readonly IConnectionMultiplexer _redis;
	private readonly ISubscriber _subscriber;
	
	public RedisSubscriber(string redisConnectionString)
	{
	_redis = ConnectionMultiplexer.Connect(redisConnectionString);
	_subscriber = _redis.GetSubscriber();
	Console.WriteLine("Redis订阅者已连接");
	}
	
	/// <summary>
	/// 订阅频道
	/// </summary>
	public async Task SubscribeAsync(string channel)
	{
	// 1. 订阅指定频道,收到消息时触发回调
	await _subscriber.SubscribeAsync(channel, (ch, msg) =>
	{
	Console.WriteLine($"收到频道 {ch} 的消息:{msg}");
	// 处理消息逻辑,比如发送短信通知用户
	});
	Console.WriteLine($"已订阅频道 {channel}");
	}
	
	public void Dispose()
	{
	_redis.Dispose();
	}
	}
	
	// 测试代码:模拟通知服务订阅订单通知
	class Program
	{
	static async Task Main(string[] args)
	{
	var subscriber = new RedisSubscriber("localhost:6379");
	await subscriber.SubscribeAsync("order_notify");
	Console.WriteLine("按任意键退出...");
	Console.ReadKey();
	subscriber.Dispose();
	}
	}

核心代码拆解:
1.连接复用器:ConnectionMultiplexer是Redis的核心,必须单例使用,避免频繁创建连接;
2.发布消息:用PublishAsync发布消息到指定频道,支持字符串、字节数组等类型;
3.订阅消息:用SubscribeAsync订阅频道,收到消息时触发回调,处理消息逻辑;
4.缺点:客户端断开后,期间的消息会丢失,不支持持久化,不能用于需要可靠消息的场景。
拓展知识:Redis Pub/Sub vs Redis Stream
特性 Redis Pub/Sub Redis Stream
消息持久化 不支持 支持
客户端断开重连 丢失消息 可以读取历史消息
消息确认 不支持 支持消费者组确认
消息堆积处理 不支持 支持
适用场景 实时通知、低可靠性 可靠消息队列、订单通知
我踩过的坑:用Redis Pub/Sub做订单通知,结果客户端断开丢失消息——换成Redis Stream后,消息持久化,客户端重连可以读取历史消息,解决了丢失问题!
三、RabbitMQ:企业级可靠消息队列,支持复杂路由与死信队列
RabbitMQ是基于AMQP协议的企业级消息队列,核心是“可靠消息+复杂路由”,支持持久化、消息确认、死信队列、多种交换机类型,适合需要可靠消息的场景(如订单通知、支付回调、异步任务)。
核心原理(大白话+快递分拣例子)
把RabbitMQ比作“快递分拣中心”:
1.生产者:寄快递的人(发送消息的服务);
2.交换机(Exchange):分拣中心,根据规则把快递分到不同的快递柜;
3.队列(Queue):快递柜,存储快递(消息);
4.绑定(Binding):分拣规则,比如“北京的快递分到A柜”;
5.消费者:取快递的人(处理消息的服务)。
我踩过的坑:没开消息确认,导致RabbitMQ重启后消息丢失——必须开启消息持久化、队列持久化、消息确认!
实战2:RabbitMQ实现可靠订单通知(C#)
步骤1:安装NuGet包
bash
Install-Package RabbitMQ.Client
步骤2:生产者代码(逐行讲解)
csharp
using RabbitMQ.Client;
using System;
using System.Text;
using System.Threading.Tasks;

namespace MessageQueue.RabbitMQ;

public class RabbitMQProducer
{
private readonly IConnection _connection;
private readonly IModel _channel;
private readonly string _exchangeName = "order_exchange";
private readonly string _routingKey = "order.notify";

public RabbitMQProducer(string rabbitMqConnectionString)
{
// 1. 创建连接工厂
var factory = new ConnectionFactory { Uri = new Uri(rabbitMqConnectionString) };
// 2. 创建连接(生产环境用单例)
_connection = factory.CreateConnection();
// 3. 创建通道
_channel = _connection.CreateModel();
// 4. 声明交换机(持久化、Direct类型)
_channel.ExchangeDeclare(_exchangeName, ExchangeType.Direct, durable: true, autoDelete: false);
// 5. 声明队列(持久化、非排他、非自动删除)
string queueName = "order_notify_queue";
_channel.QueueDeclare(queueName, durable: true, exclusive: false, autoDelete: false);
// 6. 绑定交换机和队列,指定路由键
_channel.QueueBind(queueName, _exchangeName, _routingKey);
Console.WriteLine("RabbitMQ生产者已初始化");
}

/// <summary>
/// 发送可靠消息
/// </summary>
public async Task SendMessageAsync(string message)
{
try
{
// 1. 开启消息确认(生产环境必须开启)
_channel.ConfirmSelect();
// 2. 设置消息属性:持久化、过期时间
var properties = _channel.CreateBasicProperties();
properties.Persistent = true; // 消息持久化,RabbitMQ重启后不丢失
properties.Expiration = "3600000"; // 消息1小时后过期
properties.MessageId = Guid.NewGuid().ToString("N"); // 唯一消息ID,用于幂等性

// 3. 转换消息为字节数组
byte[] body = Encoding.UTF8.GetBytes(message);
// 4. 发送消息到交换机,指定路由键
_channel.BasicPublish(_exchangeName, _routingKey, properties, body);
Console.WriteLine($"发送消息:{message},消息ID:{properties.MessageId}");

// 5. 等待消息确认,超时时间5秒
if (!await _channel.WaitForConfirmsAsync(TimeSpan.FromSeconds(5)))
{
throw new InvalidOperationException("消息未被RabbitMQ确认");
}
}
catch (Exception ex)
{
Console.WriteLine($"发送消息失败:{ex.Message}");
// 生产环境可以重试或放入死信队列
}
}

public void Dispose()
{
_channel.Close();
_connection.Close();
}
}

// 测试代码:模拟订单服务发送通知
class Program
{
static async Task Main(string[] args)
{
var producer = new RabbitMQProducer("amqp://guest:guest@localhost:5672/");
for (int i = 0; i < 10; i++)
{
string message = $"订单 {i+1} 已支付,金额:{(i+1)*100}元";
await producer.SendMessageAsync(message);
await Task.Delay(1000);
}
producer.Dispose();
}
}

核心代码拆解:
1.交换机声明:ExchangeType.Direct是直接路由,根据路由键精确匹配;
2.队列声明:durable: true表示队列持久化,RabbitMQ重启后队列不会丢失;
3.消息持久化:properties.Persistent = true表示消息持久化,RabbitMQ重启后消息不会丢失;
4.消息确认:ConfirmSelect()开启消息确认,WaitForConfirmsAsync()等待RabbitMQ确认消息已接收,避免消息丢失;
5.唯一消息ID:用于幂等性处理,避免重复消费。
步骤3:消费者代码(逐行讲解)
csharp

	using RabbitMQ.Client;
	using RabbitMQ.Client.Events;
	using System;
	using System.Text;
	using System.Threading.Tasks;
	
	namespace MessageQueue.RabbitMQ;
	
	public class RabbitMQConsumer
	{
	private readonly IConnection _connection;
	private readonly IModel _channel;
	private readonly string _queueName = "order_notify_queue";
	
	public RabbitMQConsumer(string rabbitMqConnectionString)
	{
	var factory = new ConnectionFactory { Uri = new Uri(rabbitMqConnectionString) };
	_connection = factory.CreateConnection();
	_channel = _connection.CreateModel();
	// 1. 设置预取计数,一次只处理1条消息,避免消费者过载
	_channel.BasicQos(0, 1, false);
	Console.WriteLine("RabbitMQ消费者已初始化");
	}
	
	/// <summary>
	/// 开始消费消息
	/// </summary>
	public void StartConsuming()
	{
	var consumer = new EventingBasicConsumer(_channel);
	consumer.Received += async (model, ea) =>
	{
	var body = ea.Body.ToArray();
	string message = Encoding.UTF8.GetString(body);
	string messageId = ea.BasicProperties.MessageId;
	Console.WriteLine($"收到消息:{message},消息ID:{messageId}");
	
	try
	{
	// 2. 处理消息逻辑,比如发送短信通知用户
	await ProcessMessageAsync(message);
	// 3. 手动确认消息,RabbitMQ会删除该消息
	_channel.BasicAck(ea.DeliveryTag, false);
	Console.WriteLine($"消息 {messageId} 处理成功,已确认");
	}
	catch (Exception ex)
	{
	Console.WriteLine($"消息 {messageId} 处理失败:{ex.Message}");
	// 4. 拒绝消息,重新入队(最多重试3次,否则放入死信队列)
	int retryCount = ea.BasicProperties.Headers?.ContainsKey("retry_count") == true ? (int)ea.BasicProperties.Headers["retry_count"] : 0;
	if (retryCount < 3)
	{
	retryCount++;
	ea.BasicProperties.Headers["retry_count"] = retryCount;
	_channel.BasicNack(ea.DeliveryTag, false, true); // 重新入队
	Console.WriteLine($"消息 {messageId} 重试 {retryCount} 次");
	}
	else
	{
	_channel.BasicNack(ea.DeliveryTag, false, false); // 拒绝并放入死信队列
	Console.WriteLine($"消息 {messageId} 重试3次失败,放入死信队列");
	}
	}
	};
	
	// 5. 消费队列,手动确认消息
	_channel.BasicConsume(queue: _queueName, autoAck: false, consumer: consumer);
	Console.WriteLine($"开始消费队列 {_queueName}");
	}
	
	/// <summary>
	/// 模拟处理消息逻辑
	/// </summary>
	private async Task ProcessMessageAsync(string message)
	{
	// 比如调用短信API发送通知
	await Task.Delay(100);
	// 模拟处理失败
	// if (message.Contains("订单 5")) throw new Exception("短信API调用失败");
	}
	
	public void Dispose()
	{
	_channel.Close();
	_connection.Close();
	}
	}
	
	// 测试代码:模拟通知服务消费订单通知
	class Program
	{
	static void Main(string[] args)
	{
	var consumer = new RabbitMQConsumer("amqp://guest:guest@localhost:5672/");
	consumer.StartConsuming();
	Console.WriteLine("按任意键退出...");
	Console.ReadKey();
	consumer.Dispose();
	}
	}

核心代码拆解:
1.预取计数:BasicQos(0,1,false)表示一次只处理1条消息,避免消费者过载;
2.手动确认:autoAck: false表示手动确认消息,处理成功后调用BasicAck,RabbitMQ会删除该消息;
3.重试机制:处理失败时,重试最多3次,超过则放入死信队列;
4.死信队列:需要提前声明死信交换机和死信队列,绑定到主队列,比如:
csharp
5.
6.
// 声明死信交换机
_channel.ExchangeDeclare("order_dlx_exchange", ExchangeType.Direct, durable: true);
// 声明死信队列
string dlxQueueName = "order_notify_dlx_queue";
_channel.QueueDeclare(dlxQueueName, durable: true, exclusive: false, autoDelete: false);
// 绑定死信队列到死信交换机
_channel.QueueBind(dlxQueueName, "order_dlx_exchange", "order.notify.dlx");
// 主队列绑定死信交换机
var queueArgs = new Dictionary<string, object>
{
{ "x-dead-letter-exchange", "order_dlx_exchange" },
{ "x-dead-letter-routing-key", "order.notify.dlx" }
};
_channel.QueueDeclare(_queueName, durable: true, exclusive: false, autoDelete: false, queueArgs);
7.
拓展知识:RabbitMQ四种交换机类型

类型 核心逻辑 适用场景
Direct 路由键精确匹配 订单通知、支付回调(一对一)
Fanout 广播到所有绑定的队列 实时监控、日志推送(一对多)
Topic 路由键模糊匹配(*、#) 多环境消息推送(比如prod.order.notify)
Headers 根据消息头匹配 复杂路由场景(很少用)

我踩过的坑:用Direct交换机做日志推送,结果每个环境都要创建一个队列——换成Topic交换机,用路由键prod.log、test.log,一个交换机搞定所有环境的日志推送!
四、Kafka:高吞吐量消息队列,适合大数据与日志收集
Kafka是基于分布式发布订阅的高吞吐量消息队列,核心是“高吞吐量+持久化+分区消费”,支持百万级QPS、持久化、消费者组、分区副本,适合大数据处理、日志收集、实时计算场景。
核心原理(大白话+火车运输例子)
把Kafka比作“火车运输系统”:
1.主题(Topic):火车线路,比如“北京-上海”;
2.分区(Partition):火车车厢,每个车厢存储一部分消息;
3.生产者:乘客,把行李放到指定车厢(指定分区键);
4.消费者组(Consumer Group):下客的团队,每个团队的成员负责不同的车厢(分区);
5.副本(Replica):车厢备份,避免车厢损坏导致行李丢失。
我踩过的坑:分区键选了固定值,导致所有消息都放到一个分区,QPS从10000跌到2000——换成用户ID作为分区键,消息均匀分布到多个分区,QPS恢复到10000!
实战3:Kafka实现日志收集(C#)
步骤1:安装NuGet包
bash
Install-Package Confluent.Kafka
步骤2:生产者代码(逐行讲解)
csharp
using Confluent.Kafka;
using System;
using System.Threading.Tasks;

namespace MessageQueue.Kafka;

public class KafkaProducer
{
private readonly IProducer<Null, string> _producer;
private readonly string _topicName = "app_logs";

public KafkaProducer(string kafkaBootstrapServers)
{
// 1. 配置生产者
var config = new ProducerConfig
{
BootstrapServers = kafkaBootstrapServers,
Acks = Acks.All, // 等待所有副本确认,最高可靠性
RetryBackoffMs = 1000, // 重试间隔
MessageTimeoutMs = 5000, // 消息超时时间
BatchSize = 16384, // 批量发送大小(16KB)
LingerMs = 5 // 等待5ms再发送,积累更多消息批量发送,提升吞吐量
};
// 2. 创建生产者
_producer = new ProducerBuilder<Null, string>(config).Build();
Console.WriteLine("Kafka生产者已初始化");
}

/// <summary>
/// 发送日志消息
/// </summary>
public async Task SendLogAsync(string logLevel, string message, string userId = null)
{
var logMessage = $"[{DateTime.Now:yyyy-MM-dd HH:mm:ss}] [{logLevel}] [{userId ?? "unknown"}] {message}";
// 1. 构建消息,指定分区键(用户ID),确保同一用户的消息放到同一个分区
var messageKey = userId != null ? new Message<Null, string> { Value = logMessage, Key = new Null() } : 
new Message<Null, string> { Value = logMessage, Key = new Null(), Partition = Partition.Any };
// 如果需要指定分区键,用下面的代码:
// var messageKey = new Message<string, string> { Value = logMessage, Key = userId };

try
{
// 2. 发送消息到主题
var deliveryResult = await _producer.ProduceAsync(_topicName, messageKey);
Console.WriteLine($"发送日志到分区 {deliveryResult.Partition}:{logMessage}");
}
catch (ProduceException<Null, string> ex)
{
Console.WriteLine($"发送日志失败:{ex.Error.Reason}");
}
}

public void Dispose()
{
_producer.Flush(TimeSpan.FromSeconds(10)); // 发送剩余消息
_producer.Dispose();
}
}

// 测试代码:模拟应用发送日志
class Program
{
static async Task Main(string[] args)
{
var producer = new KafkaProducer("localhost:9092");
for (int i = 0; i < 10; i++)
{
string logLevel = i % 2 == 0 ? "INFO" : "ERROR";
string message = $"用户操作日志:{i+1}";
string userId = $"user_{i%5}"; // 5个用户,消息均匀分布到5个分区
await producer.SendLogAsync(logLevel, message, userId);
await Task.Delay(100);
}
producer.Dispose();
}
}

核心代码拆解:
1.Acks配置:Acks.All表示等待所有副本确认,最高可靠性,适合需要可靠消息的场景;
2.批量发送:BatchSize和LingerMs配置批量发送,提升吞吐量;
3.分区键:指定用户ID作为分区键,确保同一用户的消息放到同一个分区,保证消息顺序;
4.消息持久化:Kafka默认持久化消息,消息会存储到磁盘,不会丢失。
步骤3:消费者代码(逐行讲解)
csharp

	using Confluent.Kafka;
	using System;
	using System.Threading;
	using System.Threading.Tasks;
	
	namespace MessageQueue.Kafka;
	
	public class KafkaConsumer
	{
	private readonly IConsumer<Null, string> _consumer;
	private readonly string _topicName = "app_logs";
	private readonly string _consumerGroup = "log_consumer_group";
	
	public KafkaConsumer(string kafkaBootstrapServers)
	{
	// 1. 配置消费者
	var config = new ConsumerConfig
	{
	BootstrapServers = kafkaBootstrapServers,
	GroupId = _consumerGroup, // 消费者组ID,同一组的消费者分摊消费分区
	AutoOffsetReset = AutoOffsetReset.Earliest, // 从头开始消费(如果没有偏移量)
	EnableAutoCommit = false, // 手动提交偏移量
	FetchMaxBytes = 52428800, // 每次拉取最大50MB
	MaxPollIntervalMs = 300000 // 最大拉取间隔(5分钟)
	};
	// 2. 创建消费者
	_consumer = new ConsumerBuilder<Null, string>(config).Build();
	// 3. 订阅主题
	_consumer.Subscribe(_topicName);
	Console.WriteLine("Kafka消费者已初始化,订阅主题:" + _topicName);
	}
	
	/// <summary>
	/// 开始消费日志
	/// </summary>
	public async Task StartConsumingAsync(CancellationToken cancellationToken)
	{
	while (!cancellationToken.IsCancellationRequested)
	{
	try
	{
	// 1. 拉取消息,超时时间1秒
	var consumeResult = _consumer.Consume(cancellationToken);
	string logMessage = consumeResult.Message.Value;
	Console.WriteLine($"收到分区 {consumeResult.Partition} 的日志:{logMessage}");
	
	// 2. 处理日志逻辑,比如存储到Elasticsearch
	await ProcessLogAsync(logMessage);
	
	// 3. 手动提交偏移量,确保消息处理成功后再提交
	_consumer.Commit(consumeResult);
	Console.WriteLine($"已提交分区 {consumeResult.Partition} 的偏移量:{consumeResult.Offset}");
	}
	catch (ConsumeException ex)
	{
	Console.WriteLine($"消费失败:{ex.Error.Reason}");
	}
	}
	}
	
	/// <summary>
	/// 模拟处理日志逻辑
	/// </summary>
	private async Task ProcessLogAsync(string logMessage)
	{
	// 比如存储到Elasticsearch
	await Task.Delay(100);
	}
	
	public void Dispose()
	{
	_consumer.Close();
	_consumer.Dispose();
	}
	}
	
	// 测试代码:模拟日志服务消费日志
	class Program
	{
	static async Task Main(string[] args)
	{
	var consumer = new KafkaConsumer("localhost:9092");
	var cancellationTokenSource = new CancellationTokenSource();
	Console.WriteLine("按任意键退出...");
	var consumeTask = consumer.StartConsumingAsync(cancellationTokenSource.Token);
	Console.ReadKey();
	cancellationTokenSource.Cancel();
	await consumeTask;
	consumer.Dispose();
	}
	}

核心代码拆解:
1.消费者组:同一消费者组的消费者分摊消费分区,比如一个主题有5个分区,同一组的5个消费者每个消费一个分区;
2.手动提交偏移量:EnableAutoCommit = false表示手动提交偏移量,处理成功后调用Commit(),避免消息丢失;
3.自动重置偏移量:AutoOffsetReset.Earliest表示如果没有偏移量,从头开始消费;
4.分区消费:每个消费者负责一个或多个分区,提升吞吐量。
拓展知识:Kafka的分区与副本机制
1.分区:主题分为多个分区,每个分区是有序的消息队列,提升吞吐量;
2.副本:每个分区有多个副本,其中一个是Leader,负责处理读写请求,其他是Follower,同步Leader的消息,Leader故障时自动切换;
3.消费者组:同一组的消费者分摊消费分区,避免重复消费;
4.偏移量:消费者记录每个分区的消费位置,重启后可以从上次的位置继续消费。
我踩过的坑:消费者组配置错误,导致多个消费者重复消费——同一消费者组的消费者分摊消费分区,不同组的消费者都会消费所有消息!
五、三种消息队列对比与生产级最佳实践

  1. 三种消息队列对比表格
特性 Redis Pub/Sub RabbitMQ Kafka
消息持久化 不支持 支持 支持
吞吐量 高 中 极高
消息确认 不支持 支持 支持
复杂路由 不支持 支持 不支持
死信队列 不支持 支持 支持
消费者组 不支持 支持 支持
适用场景 实时通知、低可靠性 企业级可靠消息、订单通知 大数据、日志收集、实时计算
  1. 生产级最佳实践
    消息幂等性:每个消息生成唯一ID,消费者处理前检查是否已处理过,避免重复消费;
    重试机制:处理失败时重试最多3次,超过则放入死信队列,避免消息堆积;
    监控告警:监控消息堆积、消费延迟、生产者/消费者错误,设置告警阈值;
    性能优化:
    RabbitMQ:开启批量发送、预取计数、持久化;
    Kafka:合理设置分区数、批量发送、消费者组数量;
    死信队列处理:定时处理死信队列的消息,比如人工审核、重新发送;
    消息过期时间:设置消息过期时间,避免无效消息堆积。
  2. 我踩过的坑总结
    1.用Redis Pub/Sub做可靠消息:客户端断开丢失消息——换成RabbitMQ或Kafka;
    2.没开消息确认:RabbitMQ重启后消息丢失——开启消息确认、持久化;
    3.分区键选择不当:Kafka所有消息放到一个分区,QPS暴跌——用用户ID、订单ID作为分区键;
    4.消费者组配置错误:多个消费者重复消费——同一组的消费者分摊消费分区;
    5.死信队列没处理:死信队列堆积大量消息——定时处理死信队列的消息。
    六、总结
    Redis Pub/Sub适合实时通知、低可靠性场景;RabbitMQ适合企业级可靠消息、复杂路由场景;Kafka适合高吞吐量、大数据、日志收集场景。选择合适的消息队列,结合生产级最佳实践,能让你的系统既可靠又高性能。
    下一节我们会学习消息队列的幂等性与重试机制,解决重复消费和消息丢失的问题。

转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49570.html


相关教程