Skip to content

消息队列重复消费 ​

概述 ​

在基于消息队列(Message Queue)的异步架构中,消息可能被重复投递和消费,原因包括:

  1. 网络故障:消费者处理成功但 ACK 丢失
  2. 消费者崩溃:处理后未及时 ACK
  3. 重试机制:消费者失败后 MQ 重新投递
  4. 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 ​

  1. 始终启用手动 ACK:不要使用自动确认
  2. 记录消息 ID:用于去重和追踪
  3. 设置合理的过期时间:避免无限增长
  4. 使用业务键去重:更可靠
  5. 监控重复率:及时发现异常
  6. 实现死信队列:处理无法消费的消息

❌ DON'T ​

  1. 不要假设消息只投递一次:MQ 至少投递一次(at-least-once)
  2. 不要忽略重复检测:这是必须的
  3. 不要在消费者中做耗时操作:先保存消息,异步处理
  4. 不要忘记清理旧数据:定期清理去重表

总结 ​

消息队列的幂等性设计是分布式系统的核心要求:

  1. 理解 MQ 语义:大多数 MQ 保证 at-least-once 投递
  2. 实现去重机制:基于消息 ID 或业务键
  3. 选择合适的存储:Redis(高性能)或数据库(强一致性)
  4. 监控和告警:及时发现重复率异常

通过合理的幂等性设计,可以安全地处理重复消息,保证数据一致性。

Released under the MIT License.