You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在ASP.NET Web API中借助RabbitMQ合并文件上传与模板创建消息

问题场景与需求

我用ASP.NET Web API + Angular开发应用,核心功能有两个:

  • 文件上传:用户上传文件成功后,服务器生成文件URL并发送到RabbitMQ队列;
  • 模板创建:用户创建模板后,把模板详情发送到另一个RabbitMQ队列。

现在需要可靠合并这两个队列的消息,最终实现:模板创建时发送带上传文件作为附件的邮件。具体疑问:

  1. 可靠合并RabbitMQ两个队列消息的最优方案是什么?
  2. 应该用专用处理服务消费两个队列消息,还是有更优策略?
  3. 如何确保只有消息成功合并后才发送邮件?

附上当前的消费者代码:

using Microsoft.Extensions.Options;
using Newtonsoft.Json;
using RabbitMQ.Client.Events;
using RabbitMQ.Client;
using System.Text;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;

public class AlertMessageConsumerService : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly RabbitMQSetting _rabbitMqSetting;
    private readonly ILogger<AlertMessageConsumerService> _logger;

    private IConnection _connection;
    private IModel _channel;

    public AlertMessageConsumerService(IOptions<RabbitMQSetting> rabbitMqSetting, IServiceProvider serviceProvider, ILogger<AlertMessageConsumerService> logger)
    {
        _rabbitMqSetting = rabbitMqSetting.Value;
        _serviceProvider = serviceProvider;
        _logger = logger;

        var factory = new ConnectionFactory
        {
            HostName = _rabbitMqSetting.HostName,
            UserName = _rabbitMqSetting.UserName,
            Password = _rabbitMqSetting.Password
        };
        _connection = factory.CreateConnection();
        _channel = _connection.CreateModel();
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        StartConsuming("alertQueue", stoppingToken);
        StartConsuming("anotherQueue", stoppingToken);

        await Task.CompletedTask;
    }

    private void StartConsuming(string queueName, CancellationToken cancellationToken)
    {
        _channel.QueueDeclare(queue: queueName, durable: false, exclusive: false, autoDelete: false, arguments: null);

        var consumer = new EventingBasicConsumer(_channel);
        consumer.Received += async (model, ea) =>
        {
            var body = ea.Body.ToArray();
            var message = Encoding.UTF8.GetString(body);

            bool processedSuccessfully = false;
            try
            {
                processedSuccessfully = await ProcessMessageAsync(message);
            }
            catch (Exception ex)
            {
                _logger.LogError($"Exception occurred while processing message from queue {queueName}: {ex}");
            }

            if (processedSuccessfully)
            {
                _channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
            }
            else
            {
                _channel.BasicReject(deliveryTag: ea.DeliveryTag, requeue: true);
            }
        };

        _channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
    }

    private async Task<bool> ProcessMessageAsync(string message)
    {
        try
        {
            using (var scope = _serviceProvider.CreateScope())
            {
                var emailService = scope.ServiceProvider.GetRequiredService<IEmailSender>();
                var vsoService = scope.ServiceProvider.GetRequiredService<IVsoService>();
                var alertTypeService = scope.ServiceProvider.GetRequiredService<IAlertTypeService>();
                var s3Service = scope.ServiceProvider.GetRequiredService<IFileService>();

                var alertMessage = JsonConvert.DeserializeObject<AlertMessage>(message);
                var uploadMessage = JsonConvert.DeserializeObject<FileUploadMessage>(message);

                if (alertMessage != null && uploadMessage != null)
                {
                    _logger.LogInformation($"AlertMessage: {JsonConvert.SerializeObject(alertMessage)}");
                    _logger.LogInformation($"FileUploadMessage: {JsonConvert.SerializeObject(uploadMessage)}");

                    return true;
                }
                else
                {
                    _logger.LogWarning("Both AlertMessage and FileUploadMessage must be present in the message.");
                    return false;
                }
            }
        }
        catch (JsonException jsonEx)
        {
            _logger.LogError($"JSON error processing message: {jsonEx.Message}");
            return false;
        }
        catch (Exception ex)
        {
            _logger.LogError($"Error processing message: {ex.Message}");
            return false;
        }
    }

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

解决方案与建议

1. 可靠合并RabbitMQ两个队列消息的最优方案

最优方案是基于关联ID的消息暂存+状态跟踪,核心逻辑:

  • 给每一组需要合并的消息分配唯一CorrelationId:可以在用户发起模板创建请求时生成,上传文件和创建模板的流程全程携带这个ID,发送RabbitMQ消息时将其作为消息属性或体字段。
  • 用持久化存储(Redis/数据库)暂存单条消息:收到文件URL消息时,以CorrelationId为键存储;收到模板详情消息时,用同一ID查找是否已有对应文件消息,反之亦然。
  • 同一CorrelationId的两条消息到齐后,触发合并处理。

可靠性补充:

  • 给RabbitMQ队列设置durable: true,消息标记persistent: true,避免服务重启丢失消息。
  • 给暂存消息设置过期时间(比如1小时),超时则触发告警或清理,防止因单条消息丢失导致永久等待。

2. 是否使用专用处理服务?

必须用专用的消息合并处理服务,理由:

  • 职责单一:专门负责消费队列、暂存消息、合并逻辑、触发邮件发送,避免与业务API耦合。
  • 可扩展性:后续新增合并需求时,只需扩展该服务逻辑,无需修改业务代码。
  • 可靠性集中控制:可以在服务内统一处理消息确认、重试、死信队列等逻辑,避免业务API被阻塞。

推荐采用ASP.NET Core BackgroundService的形式(和你当前代码结构一致),规模小的话可以和业务服务同进程部署,规模大则建议独立成微服务。

3. 确保仅合并成功后发送邮件的措施

要实现这个目标,需在合并和邮件环节做以下控制:

  • 原子性处理:两条消息到齐后,先执行邮件发送逻辑,只有发送成功后,再:
    1. 从暂存存储中删除对应消息记录;
    2. 向RabbitMQ发送BasicAck确认消息已处理。
  • 失败重试机制:邮件发送失败时,不要确认消息,将其重新入队(或送入死信队列),同时保留暂存的消息记录,等待重试。
  • 幂等性保障:用CorrelationId作为邮件发送的唯一标识,记录已发送的邮件ID,避免重复发送。

优化后的代码示例

以下是调整后的消费者服务,实现了基于Redis的消息暂存与合并逻辑:

using Microsoft.Extensions.Options;
using Newtonsoft.Json;
using RabbitMQ.Client.Events;
using RabbitMQ.Client;
using System.Text;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using StackExchange.Redis;

public class MessageMergeConsumerService : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly RabbitMQSetting _rabbitMqSetting;
    private readonly ILogger<MessageMergeConsumerService> _logger;
    private readonly IDatabase _redisDb;

    private IConnection _connection;
    private IModel _channel;
    private const string FileMessagePrefix = "file_";
    private const string TemplateMessagePrefix = "template_";
    private const int MessageExpirySeconds = 3600; // 1小时过期

    public MessageMergeConsumerService(IOptions<RabbitMQSetting> rabbitMqSetting, 
                                       IServiceProvider serviceProvider, 
                                       ILogger<MessageMergeConsumerService> logger,
                                       IConnectionMultiplexer redisMultiplexer)
    {
        _rabbitMqSetting = rabbitMqSetting.Value;
        _serviceProvider = serviceProvider;
        _logger = logger;
        _redisDb = redisMultiplexer.GetDatabase();

        var factory = new ConnectionFactory
        {
            HostName = _rabbitMqSetting.HostName,
            UserName = _rabbitMqSetting.UserName,
            Password = _rabbitMqSetting.Password,
            DispatchConsumersAsync = true // 支持异步消费者
        };
        _connection = factory.CreateConnection();
        _channel = _connection.CreateModel();
        // 声明持久化队列
        _channel.QueueDeclare(queue: "fileUploadQueue", durable: true, exclusive: false, autoDelete: false, arguments: null);
        _channel.QueueDeclare(queue: "templateCreateQueue", durable: true, exclusive: false, autoDelete: false, arguments: null);
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        await StartConsumingAsync("fileUploadQueue", ProcessFileMessageAsync, stoppingToken);
        await StartConsumingAsync("templateCreateQueue", ProcessTemplateMessageAsync, stoppingToken);
        await Task.CompletedTask;
    }

    private async Task StartConsumingAsync(string queueName, Func<string, BasicDeliverEventArgs, Task<bool>> processFunc, CancellationToken cancellationToken)
    {
        var consumer = new AsyncEventingBasicConsumer(_channel);
        consumer.Received += async (model, ea) =>
        {
            var body = ea.Body.ToArray();
            var message = Encoding.UTF8.GetString(body);
            bool processedSuccessfully = false;

            try
            {
                processedSuccessfully = await processFunc(message, ea);
            }
            catch (Exception ex)
            {
                _logger.LogError($"处理队列{queueName}消息出错: {ex}");
                // 重试3次后送入死信队列
                if (ea.BasicProperties.Headers.TryGetValue("x-retry-count", out var retryCountObj) && 
                    int.TryParse(retryCountObj.ToString(), out int retryCount) && retryCount >=3)
                {
                    _channel.BasicReject(ea.DeliveryTag, false);
                }
                else
                {
                    var props = _channel.CreateBasicProperties();
                    props.Headers = ea.BasicProperties.Headers ?? new Dictionary<string, object>();
                    props.Headers["x-retry-count"] = (ea.BasicProperties.Headers?.ContainsKey("x-retry-count") ?? false) 
                        ? (int)ea.BasicProperties.Headers["x-retry-count"] + 1 
                        : 1;
                    _channel.BasicPublish(ea.Exchange, ea.RoutingKey, props, ea.Body);
                    _channel.BasicAck(ea.DeliveryTag, false);
                }
                return;
            }

            if (processedSuccessfully)
            {
                _channel.BasicAck(ea.DeliveryTag, false);
            }
            else
            {
                _channel.BasicReject(ea.DeliveryTag, true); // 重新入队
            }
        };

        _channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
        _logger.LogInformation($"开始消费队列{queueName}");
    }

    private async Task<bool> ProcessFileMessageAsync(string message, BasicDeliverEventArgs ea)
    {
        try
        {
            var fileMsg = JsonConvert.DeserializeObject<FileUploadMessage>(message);
            if (fileMsg == null || string.IsNullOrEmpty(fileMsg.CorrelationId))
            {
                _logger.LogWarning("文件消息格式无效或缺少CorrelationId");
                return false;
            }

            // 暂存文件消息到Redis
            await _redisDb.StringSetAsync($"{FileMessagePrefix}{fileMsg.CorrelationId}", message, TimeSpan.FromSeconds(MessageExpirySeconds));
            
            // 检查对应模板消息是否存在
            var templateMsgStr = await _redisDb.StringGetAsync($"{TemplateMessagePrefix}{fileMsg.CorrelationId}");
            if (!templateMsgStr.IsNullOrEmpty)
            {
                return await ProcessMergedMessagesAsync(fileMsg, JsonConvert.DeserializeObject<TemplateCreateMessage>(templateMsgStr));
            }

            _logger.LogInformation($"暂存文件消息,CorrelationId: {fileMsg.CorrelationId}");
            return true;
        }
        catch (JsonException jsonEx)
        {
            _logger.LogError($"解析文件消息JSON出错: {jsonEx.Message}");
            return false;
        }
    }

    private async Task<bool> ProcessTemplateMessageAsync(string message, BasicDeliverEventArgs ea)
    {
        try
        {
            var templateMsg = JsonConvert.DeserializeObject<TemplateCreateMessage>(message);
            if (templateMsg == null || string.IsNullOrEmpty(templateMsg.CorrelationId))
            {
                _logger.LogWarning("模板消息格式无效或缺少CorrelationId");
                return false;
            }

            // 暂存模板消息到Redis
            await _redisDb.StringSetAsync($"{TemplateMessagePrefix}{templateMsg.CorrelationId}", message, TimeSpan.FromSeconds(MessageExpirySeconds));
            
            // 检查对应文件消息是否存在
            var fileMsgStr = await _redisDb.StringGetAsync($"{FileMessagePrefix}{templateMsg.CorrelationId}");
            if (!fileMsgStr.IsNullOrEmpty)
            {
                return await ProcessMergedMessagesAsync(JsonConvert.DeserializeObject<FileUploadMessage>(fileMsgStr), templateMsg);
            }

            _logger.LogInformation($"暂存模板消息,CorrelationId: {templateMsg.CorrelationId}");
            return true;
        }
        catch (JsonException jsonEx)
        {
            _logger.LogError($"解析模板消息JSON出错: {jsonEx.Message}");
            return false;
        }
    }

    private async Task<bool> ProcessMergedMessagesAsync(FileUploadMessage fileMsg, TemplateCreateMessage templateMsg)
    {
        using var scope = _serviceProvider.CreateScope();
        var emailService = scope.ServiceProvider.GetRequiredService<IEmailSender>();
        var fileService = scope.ServiceProvider.GetRequiredService<IFileService>();

        try
        {
            // 下载文件内容
            var fileContent = await fileService.DownloadFileAsync(fileMsg.FileUrl);
            // 发送带附件的邮件
            await emailService.SendEmailWithAttachmentAsync(
                templateMsg.RecipientEmail,
                templateMsg.EmailSubject,
                templateMsg.EmailContent,
                fileMsg.FileName,
                fileContent);

            // 发送成功后清理Redis暂存
            await _redisDb.KeyDeleteAsync($"{FileMessagePrefix}{fileMsg.CorrelationId}");
            await _redisDb.KeyDeleteAsync($"{TemplateMessagePrefix}{fileMsg.CorrelationId}");

            _logger.LogInformation($"合并消息并发送邮件成功,CorrelationId: {fileMsg.CorrelationId}");
            return true;
        }
        catch (Exception ex)
        {
            _logger.LogError($"合并消息处理出错,CorrelationId: {fileMsg.CorrelationId},错误: {ex}");
            return false;
        }
    }

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

// 消息模型示例
public class FileUploadMessage
{
    public string CorrelationId { get; set; }
    public string FileUrl { get; set; }
    public string FileName { get; set; }
}

public class TemplateCreateMessage
{
    public string CorrelationId { get; set; }
    public string RecipientEmail { get; set; }
    public string EmailSubject { get; set; }
    public string EmailContent { get; set; }
}

关键注意事项
  • 确保CorrelationId在文件上传和模板创建流程中全程传递,可由前端发起请求时生成或后端统一分配。
  • Redis需配置持久化,避免重启丢失暂存消息。
  • 给RabbitMQ配置死信队列,处理重试多次仍失败的消息,防止队列阻塞。
  • 邮件发送逻辑需添加幂等校验,比如用CorrelationId记录已发送的邮件,避免重复发送。

内容的提问来源于stack exchange,提问作者PASHA

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 15:20:54