Rebus发送Direct消息报错:无法在已完成事务上下文添加OnCommit操作
问题描述
错误信息
Error: InvalidOperationException: Cannot add OnCommit action on a completed transaction context.
Rebus.Exceptions.RebusApplicationException: 'Could not 'GetOrAdd' item with key 'outgoing-messages' as type System.Collections.Concurrent.ConcurrentQueue`1[Rebus.Transport.OutgoingTransportMessage]'
触发代码
await retryPolicy.ExecuteAsync(async () => { if (messageType == MessageType.Direct) // 调用此行时报错 await _messageBus.Send(message); else await _messageBus.Publish(message); });
场景说明
- 微服务部署在Azure和AKS环境
- MVC项目向Mail微服务发送Direct邮件消息正常运行
- Mail微服务处理邮件后,向Audit微服务发送Direct审计消息时触发上述错误
- 使用PubSub类型队列无异常,但业务要求同时支持Direct和PubSub队列
相关配置与代码
1. MVC服务Rebus配置(发送邮件消息)
var messageQueuesConfig = configuration.GetRequiredSection("MessageQueues"); var transport = messageQueuesConfig.GetValue<MessageBusTransportType>("Transport"); var hostName = messageQueuesConfig.GetValue<string>("HostName"); var userName = messageQueuesConfig.GetValue<string>("UserName"); var password = messageQueuesConfig.GetValue<string>("Password"); var queuePrefix = messageQueuesConfig.GetValue<string>("QueuePrefix"); var connectionString = MessageBusConnectionStringFactory.GetConnectionString(transport, hostName, userName, password); services.AddRebus((config, provider) => config .Transport(t => t.UseRabbitMq(connectionString, MessageBusQueueNameFactory.GetQueueName(MessageBusQueueIdentifier.Web, queuePrefix, MessageBusQueueType.PubSub))) .Options(o => o.RetryStrategy(secondLevelRetriesEnabled: true)) .Routing(r => r.TypeBased() .Map<AuditEvent>(MessageBusQueueNameFactory.GetQueueName(MessageBusQueueIdentifier.Audit, queuePrefix, MessageBusQueueType.Direct)) .Map<MailMessage>(MessageBusQueueNameFactory.GetQueueName(MessageBusQueueIdentifier.Mail, queuePrefix, MessageBusQueueType.Direct)) ), onCreated: async bus => { await bus.Subscribe<Event>(); } );
2. Mail微服务Rebus配置(处理邮件消息)
var messageQueuesConfig = configuration.GetRequiredSection("MessageQueues"); var transport = messageQueuesConfig.GetValue<MessageBusTransportType>("Transport"); var hostName = messageQueuesConfig.GetValue<string>("HostName"); var userName = messageQueuesConfig.GetValue<string>("UserName"); var password = messageQueuesConfig.GetValue<string>("Password"); var queuePrefix = messageQueuesConfig.GetValue<string>("QueuePrefix"); var connectionString = MessageBusConnectionStringFactory.GetConnectionString(transport, hostName, userName, password); services.AddRebus((config, provider) => config .Transport(t => t.UseRabbitMq(connectionString, MessageBusQueueNameFactory.GetQueueName(MessageBusQueueIdentifier.Mail, queuePrefix, MessageBusQueueType.Direct))) .Options(o => o.RetryStrategy(secondLevelRetriesEnabled: true)) .Routing(r => r.TypeBased() .Map<AuditEvent>(MessageBusQueueNameFactory.GetQueueName(MessageBusQueueIdentifier.Audit, queuePrefix, MessageBusQueueType.Direct)) ), onCreated: async bus => { await bus.Subscribe<Event>(); } );
3. Mail微服务处理邮件并发送审计消息的代码
private readonly IBus _bus; private readonly IMailClientFactory _mailClientFactory; private readonly IAuditor _auditor; private readonly ILogger<MailMessageHandler> _logger; private readonly NotificationOptions _notificationOptions; private readonly SemaphoreSlim _semaphore; public MailMessageHandler( IBus bus, IMailClientFactory mailClientFactory, IAuditor auditor, ILogger<MailMessageHandler> logger, IOptionsSnapshot<NotificationOptions> notificationOptions ) { _bus = bus; _mailClientFactory = mailClientFactory; _auditor = auditor; _logger = logger; if (notificationOptions == null) throw new ArgumentNullException(nameof(notificationOptions)); _notificationOptions = notificationOptions.Value; _semaphore = new SemaphoreSlim(1); } public async Task Handle(MailMessage message) { // 邮件处理的验证与初始化 try { _semaphore.Wait(); await retryPolicy.ExecuteAsync(async () => { var result = await mailClient.SendEmailAsync(message); if (result.Errors.Any()) { // 记录错误 } else { _ = Task.Run(() => { var auditEvent = new AuditEvent { AuditEventCode = AuditEventCode.EmailSent, EventData = System.Text.Json.JsonSerializer.Serialize(new { // 自定义数据 }) }; // 调用此行后触发问题 _auditor.AuditEvent(auditEvent); }); } }); } catch (Exception ex) { // 异常处理逻辑 } finally { _semaphore.Release(); } }
Auditor类在Mail微服务中注册为Singleton,其AuditEvent方法代码:
public void AuditEvent(AuditEvent auditEvent) { ArgumentNullException.ThrowIfNull(auditEvent); if (auditEvent.AuditEventCode == AuditEventCode.NotSet) throw new InvalidOperationException("AuditEventCode must be set before calling this method."); _messageBus.SendDirectMessage(auditEvent); }
问题定位与解决方案
核心原因
错误本质是Rebus的事务上下文已完成,但仍尝试在该上下文中发送消息。具体触发逻辑:
- Mail服务的消息处理流程中,使用
Task.Run开启了脱离当前Rebus事务上下文的后台任务 - 原处理方法的
await retryPolicy.ExecuteAsync执行完成后,Rebus自动提交并结束当前事务上下文 - 后台任务中调用
_auditor.AuditEvent发送Direct消息时,已无有效事务上下文可用,导致报错
解决方案
1. 移除Task.Run,在事务上下文内完成消息发送
将后台任务改为直接异步执行,确保消息发送操作处于原事务上下文生命周期内:
else { var auditEvent = new AuditEvent { AuditEventCode = AuditEventCode.EmailSent, EventData = System.Text.Json.JsonSerializer.Serialize(new { // 自定义数据 }) }; await _auditor.AuditEventAsync(auditEvent); // 修改Auditor为异步方法 }
同时修改Auditor的方法为异步:
public async Task AuditEventAsync(AuditEvent auditEvent) { ArgumentNullException.ThrowIfNull(auditEvent); if (auditEvent.AuditEventCode == AuditEventCode.NotSet) throw new InvalidOperationException("AuditEventCode must be set before calling this method."); await _messageBus.SendDirectMessageAsync(auditEvent); // 对应异步发送方法 }
2. 若需异步解耦,手动创建独立事务上下文发送
如果必须将审计操作解耦,可手动创建新的Rebus事务上下文发送消息:
// 在Auditor中注入IRebusTransactionScopeFactory private readonly IRebusTransactionScopeFactory _transactionScopeFactory; public async Task AuditEventAsync(AuditEvent auditEvent) { ArgumentNullException.ThrowIfNull(auditEvent); if (auditEvent.AuditEventCode == AuditEventCode.NotSet) throw new InvalidOperationException("AuditEventCode must be set before calling this method."); using var scope = _transactionScopeFactory.Create(); await _messageBus.SendDirectMessageAsync(auditEvent); await scope.CompleteAsync(); }
3. 确认Singleton Auditor的依赖合法性
Auditor注册为Singleton时,需确保其注入的_messageBus是Rebus官方提供的IBus实例(Rebus的IBus本身支持单例安全访问),避免自定义封装的_messageBus存在上下文绑定问题。
内容的提问来源于stack exchange,提问作者Brad Pickering-Dunn
相关产品推荐
相关产品推荐

