在ASP.NET Web API中借助RabbitMQ合并文件上传与模板创建消息
问题场景与需求
我用ASP.NET Web API + Angular开发应用,核心功能有两个:
- 文件上传:用户上传文件成功后,服务器生成文件URL并发送到RabbitMQ队列;
- 模板创建:用户创建模板后,把模板详情发送到另一个RabbitMQ队列。
现在需要可靠合并这两个队列的消息,最终实现:模板创建时发送带上传文件作为附件的邮件。具体疑问:
- 可靠合并RabbitMQ两个队列消息的最优方案是什么?
- 应该用专用处理服务消费两个队列消息,还是有更优策略?
- 如何确保只有消息成功合并后才发送邮件?
附上当前的消费者代码:
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. 确保仅合并成功后发送邮件的措施
要实现这个目标,需在合并和邮件环节做以下控制:
- 原子性处理:两条消息到齐后,先执行邮件发送逻辑,只有发送成功后,再:
- 从暂存存储中删除对应消息记录;
- 向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
相关产品推荐
相关产品推荐

