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

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的事务上下文已完成,但仍尝试在该上下文中发送消息。具体触发逻辑:

  1. Mail服务的消息处理流程中,使用Task.Run开启了脱离当前Rebus事务上下文的后台任务
  2. 原处理方法的await retryPolicy.ExecuteAsync执行完成后,Rebus自动提交并结束当前事务上下文
  3. 后台任务中调用_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:37:03