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

C# MassTransit运行时创建队列后消费消息抛ObjectDisposedException

问题分析与解决方案

问题核心

运行时动态创建队列后发送消息,抛出System.ObjectDisposedException,重启应用后恢复正常。

原因

代码中使用using var provider = this.serviceProvider.CreateScope();创建了临时依赖注入作用域,Consume方法执行完毕后该作用域会被自动释放。而通过c.Consumer(() => provider.ServiceProvider.GetRequiredService<UpdateMachineStatusConsumer>())获取的消费者实例,其生命周期绑定到这个临时作用域上。当队列收到消息需要处理时,消费者实例已经被销毁,因此抛出对象已释放的异常。

重启后正常是因为启动阶段的后台服务创建队列时,使用的是应用级别的长期服务提供者,消费者实例不会被提前释放。

解决方案

方案1:让MassTransit自动管理消费者生命周期

移除临时作用域,直接让MassTransit使用应用根服务提供者解析消费者,确保消费者生命周期与队列接收端点一致:

public async Task Consume(ConsumeContext<MachineCreatedMessage> context)
{
    var message = context.Message;
    var machineId = message.MachineId;
    var cancellationToken = context.CancellationToken;

    var createMachineStatusQueue = this.bus.ConnectReceiveEndpoint(
        string.Format("update_machine_status__{0}", machineId),
        x =>
        {
            if (x is IRabbitMqReceiveEndpointConfigurator c)
            {
                c.ConfigureConsumeTopology = false;
                c.ConcurrentMessageLimit = 1;
                c.PrefetchCount = 1;

                // 让MassTransit使用根服务提供者解析消费者,自动处理生命周期
                c.Consumer<UpdateMachineStatusConsumer>(this.serviceProvider);

                c.Bind("machine_status_cycle", s =>
                {
                    s.RoutingKey = machineId.ToString();
                    s.ExchangeType = ExchangeType.Direct;
                });
            }
        });

    var anotherQueues = ...
}

方案2:保留作用域并长期持有(适用于需要独立作用域的场景)

如果消费者必须使用独立作用域,不要用using自动释放,而是将作用域存储到一个长期存在的容器中,直到队列不再使用时手动释放:

// 类级别声明集合存储作用域,避免自动释放
private readonly List<IServiceScope> _scopes = new();

public async Task Consume(ConsumeContext<MachineCreatedMessage> context)
{
    var message = context.Message;
    var machineId = message.MachineId;
    var cancellationToken = context.CancellationToken;

    var provider = this.serviceProvider.CreateScope();
    _scopes.Add(provider);

    var createMachineStatusQueue = this.bus.ConnectReceiveEndpoint(
        string.Format("update_machine_status__{0}", machineId),
        x =>
        {
            if (x is IRabbitMqReceiveEndpointConfigurator c)
            {
                c.ConfigureConsumeTopology = false;
                c.ConcurrentMessageLimit = 1;
                c.PrefetchCount = 1;

                c.Consumer(() => 
                    provider.ServiceProvider.GetRequiredService<UpdateMachineStatusConsumer>()
                );

                c.Bind("machine_status_cycle", s =>
                {
                    s.RoutingKey = machineId.ToString();
                    s.ExchangeType = ExchangeType.Direct;
                });

                // 端点停止时释放对应作用域,避免内存泄漏
                c.StopCallback = async () => 
                {
                    provider.Dispose();
                    _scopes.Remove(provider);
                };
            }
        });

    var anotherQueues = ...
}

注意事项

  • 若UpdateMachineStatusConsumer依赖Scoped服务,推荐使用方案1,MassTransit会自动为每个消息处理创建新的作用域(需确保MassTransit已配置支持Scoped消费者)。
  • 长期持有作用域时,必须在队列不再使用时手动释放,避免内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 12:32:13