消息队列重复消费
概述
在基于消息队列(Message Queue)的异步架构中,消息可能被重复投递和消费,原因包括:
- 网络故障:消费者处理成功但 ACK 丢失
- 消费者崩溃:处理后未及时 ACK
- 重试机制:消费者失败后 MQ 重新投递
- MQ 重启:未持久化的 ACK 信息丢失
为什么需要幂等性
生产者 消息队列 消费者
| | |
|-- 发送消息 A ----------->| |
| |-- 投递消息 A ---------->|
| | |-- 处理消息 A ✓
| | |-- 准备 ACK
| | | (崩溃!) ❌
| | |
| |-- 重新投递消息 A ------>|
| | |-- 再次处理 A ❌
| | |-- 重复数据!如果没有幂等性保证,重复消费会导致:
- 订单重复创建
- 账户余额重复增加
- 库存重复扣减
- 通知重复发送
RabbitMQ 实现
1. 基础消费者(非幂等 - 危险)
csharp
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
public class OrderConsumer : BackgroundService
{
private readonly IConnection _connection;
private readonly IModel _channel;
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
var consumer = new EventingBasicConsumer(_channel);
consumer.Received += async (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
var order = JsonSerializer.Deserialize<OrderCreatedEvent>(message);
// ❌ 危险:如果重复消费,会创建重复订单
await _orderService.CreateOrderAsync(order);
// 手动 ACK
_channel.BasicAck(ea.DeliveryTag, false);
};
_channel.BasicConsume(queue: "orders", autoAck: false, consumer: consumer);
return Task.CompletedTask;
}
}2. 基于消息 ID 的幂等性(推荐)
csharp
public class IdempotentOrderConsumer : BackgroundService
{
private readonly IConnection _connection;
private readonly IModel _channel;
private readonly IProcessedMessagesStore _processedStore;
private readonly ILogger<IdempotentOrderConsumer> _logger;
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
var consumer = new EventingBasicConsumer(_channel);
consumer.Received += async (model, ea) =>
{
var messageId = ea.BasicProperties.MessageId;
try
{
// 检查是否已处理
if (await _processedStore.ExistsAsync(messageId))
{
_logger.LogInformation("Message {MessageId} already processed, skipping", messageId);
_channel.BasicAck(ea.DeliveryTag, false);
return;
}
// 处理消息
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
var order = JsonSerializer.Deserialize<OrderCreatedEvent>(message);
await _orderService.CreateOrderAsync(order);
// 记录已处理
await _processedStore.AddAsync(messageId, TimeSpan.FromDays(7));
// ACK
_channel.BasicAck(ea.DeliveryTag, false);
_logger.LogInformation("Message {MessageId} processed successfully", messageId);
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to process message {MessageId}", messageId);
// NACK 并重新入队(或进入死信队列)
_channel.BasicNack(ea.DeliveryTag, false, true);
}
};
_channel.BasicConsume(queue: "orders", autoAck: false, consumer: consumer);
return Task.CompletedTask;
}
}
// 已处理消息存储
public interface IProcessedMessagesStore
{
Task<bool> ExistsAsync(string messageId);
Task AddAsync(string messageId, TimeSpan expiry);
}
public class RedisProcessedMessagesStore : IProcessedMessagesStore
{
private readonly IDatabase _redis;
private const string Prefix = "mq:processed:";
public async Task<bool> ExistsAsync(string messageId)
{
return await _redis.KeyExistsAsync($"{Prefix}{messageId}");
}
public async Task AddAsync(string messageId, TimeSpan expiry)
{
await _redis.StringSetAsync(
$"{Prefix}{messageId}",
DateTime.UtcNow.ToString("O"),
expiry);
}
}3. 生产者设置消息 ID
csharp
public class OrderProducer
{
private readonly IModel _channel;
public void PublishOrderCreated(OrderCreatedEvent orderEvent)
{
var message = JsonSerializer.Serialize(orderEvent);
var body = Encoding.UTF8.GetBytes(message);
var properties = _channel.CreateBasicProperties();
properties.Persistent = true; // 持久化
properties.MessageId = Guid.NewGuid().ToString(); // 唯一消息 ID
properties.Timestamp = new AmqpTimestamp(DateTimeOffset.UtcNow.ToUnixTimeSeconds());
_channel.BasicPublish(
exchange: "orders",
routingKey: "order.created",
basicProperties: properties,
body: body);
}
}Kafka 实现
1. 利用 Kafka 的特性
Kafka 提供了更强的幂等性支持:
csharp
using Confluent.Kafka;
public class KafkaOrderConsumer : BackgroundService
{
private readonly IConsumer<string, string> _consumer;
private readonly IProcessedMessagesStore _processedStore;
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
_consumer.Subscribe("orders");
while (!stoppingToken.IsCancellationRequested)
{
try
{
var consumeResult = _consumer.Consume(stoppingToken);
var messageId = consumeResult.Message.Headers
.FirstOrDefault(h => h.Key == "message-id")
.Value != null
? Encoding.UTF8.GetString(consumeResult.Message.Headers
.First(h => h.Key == "message-id").Value)
: null;
if (string.IsNullOrEmpty(messageId))
{
_logger.LogWarning("Message without ID received");
continue;
}
// 检查是否已处理
if (await _processedStore.ExistsAsync(messageId))
{
_logger.LogInformation("Duplicate message {MessageId}, skipping", messageId);
continue;
}
// 处理消息
var order = JsonSerializer.Deserialize<OrderCreatedEvent>(
consumeResult.Message.Value);
await _orderService.CreateOrderAsync(order);
// 记录已处理
await _processedStore.AddAsync(messageId, TimeSpan.FromDays(7));
// 提交偏移量
_consumer.Commit(consumeResult);
}
catch (ConsumeException ex)
{
_logger.LogError(ex, "Consume error");
}
}
}
}2. Kafka Exactly-Once 语义
Kafka 0.11+ 支持事务性生产者:
csharp
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092",
EnableIdempotence = true, // 启用幂等性
Acks = Acks.All, // 等待所有副本确认
MaxInFlight = 1 // 保证消息顺序
};
using var producer = new ProducerBuilder<string, string>(config).Build();
// 事务性发送
producer.InitTransactions();
producer.BeginTransaction();
try
{
producer.Produce("orders", new Message<string, string>
{
Key = order.UserId.ToString(),
Value = JsonSerializer.Serialize(order),
Headers = new Headers
{
{ "message-id", Guid.NewGuid().ToByteArray() }
}
});
producer.CommitTransaction();
}
catch
{
producer.AbortTransaction();
throw;
}数据库去重表方案
1. 创建去重表
sql
-- PostgreSQL
CREATE TABLE consumed_messages (
message_id VARCHAR(128) PRIMARY KEY,
source VARCHAR(50) NOT NULL, -- rabbitmq, kafka, etc.
topic VARCHAR(100) NOT NULL,
consumed_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
processed_data JSONB, -- 可选:存储处理结果
-- 自动清理:保留 7 天
expires_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() + INTERVAL '7 days'
);
CREATE INDEX idx_consumed_messages_expires_at
ON consumed_messages(expires_at);
-- 定期清理函数
CREATE OR REPLACE FUNCTION cleanup_expired_messages()
RETURNS void AS $$
BEGIN
DELETE FROM consumed_messages
WHERE expires_at < NOW();
END;
$$ LANGUAGE plpgsql;2. C# 实现
csharp
public class DatabaseProcessedMessagesStore : IProcessedMessagesStore
{
private readonly AppDbContext _dbContext;
public async Task<bool> ExistsAsync(string messageId)
{
return await _dbContext.ConsumedMessages
.AnyAsync(m => m.MessageId == messageId);
}
public async Task AddAsync(string messageId, TimeSpan expiry)
{
var record = new ConsumedMessage
{
MessageId = messageId,
Source = "rabbitmq",
Topic = "orders",
ConsumedAt = DateTime.UtcNow,
ExpiresAt = DateTime.UtcNow.Add(expiry)
};
try
{
_dbContext.ConsumedMessages.Add(record);
await _dbContext.SaveChangesAsync();
}
catch (DbUpdateException ex) when (IsUniqueViolation(ex))
{
// 并发插入,消息已被其他实例处理
_logger.LogWarning("Concurrent consumption of message {MessageId}", messageId);
}
}
}业务级幂等性
1. 基于业务键的去重
csharp
public class PaymentConsumer
{
public async Task ConsumePaymentEvent(PaymentEvent paymentEvent)
{
// 不依赖消息 ID,而是使用业务键
var businessKey = $"payment_{paymentEvent.OrderId}_{paymentEvent.TransactionId}";
// 检查业务层面是否已处理
var existingPayment = await _payments
.FirstOrDefaultAsync(p =>
p.OrderId == paymentEvent.OrderId &&
p.TransactionId == paymentEvent.TransactionId);
if (existingPayment != null)
{
_logger.LogInformation("Payment already processed for order {OrderId}",
paymentEvent.OrderId);
return;
}
// 处理支付
var payment = new Payment
{
OrderId = paymentEvent.OrderId,
TransactionId = paymentEvent.TransactionId,
Amount = paymentEvent.Amount,
Status = "completed"
};
_payments.Add(payment);
await _dbContext.SaveChangesAsync();
}
}优势:
- 即使消息 ID 不同(如重试),业务层面也能识别重复
- 更符合业务语义
- 便于人工排查问题
监控与告警
1. 指标收集
csharp
public class ConsumerMetrics
{
private readonly Counter<long> _messagesConsumed;
private readonly Counter<long> _duplicateMessages;
private readonly Histogram<double> _processingTime;
public void RecordMessageConsumed(bool isDuplicate, double processingTimeMs)
{
_messagesConsumed.Add(1);
if (isDuplicate)
{
_duplicateMessages.Add(1);
}
_processingTime.Record(processingTimeMs);
}
}2. Prometheus 告警
yaml
groups:
- name: mq_consumer_alerts
rules:
- alert: HighDuplicateRate
expr: rate(mq_duplicate_messages_total[5m]) / rate(mq_messages_consumed_total[5m]) > 0.1
for: 5m
annotations:
summary: "High duplicate message rate"
description: "More than 10% of messages are duplicates"
- alert: ConsumerLag
expr: kafka_consumer_group_lag > 10000
for: 5m
annotations:
summary: "High consumer lag"最佳实践总结
✅ DO
- 始终启用手动 ACK:不要使用自动确认
- 记录消息 ID:用于去重和追踪
- 设置合理的过期时间:避免无限增长
- 使用业务键去重:更可靠
- 监控重复率:及时发现异常
- 实现死信队列:处理无法消费的消息
❌ DON'T
- 不要假设消息只投递一次:MQ 至少投递一次(at-least-once)
- 不要忽略重复检测:这是必须的
- 不要在消费者中做耗时操作:先保存消息,异步处理
- 不要忘记清理旧数据:定期清理去重表
总结
消息队列的幂等性设计是分布式系统的核心要求:
- 理解 MQ 语义:大多数 MQ 保证 at-least-once 投递
- 实现去重机制:基于消息 ID 或业务键
- 选择合适的存储:Redis(高性能)或数据库(强一致性)
- 监控和告警:及时发现重复率异常
通过合理的幂等性设计,可以安全地处理重复消息,保证数据一致性。