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

MassTransit动态添加消费者时出现配置异常问题排查

.NET 6 + MassTransit 8.0.5 + RabbitMQ 配置异常问题解决

我在.NET 6微服务中使用MassTransit 8.0.5搭配RabbitMQ实现服务总线,Service-A发布消息后能看到按命名空间创建的交换器,但没有队列。启动消费者Service-B时抛出如下配置异常:

配置异常截图

原错误配置代码

public static IServiceCollection AddMassTransit(this IServiceCollection services, Assembly assembly)
{
    var serviceProvider = services.BuildServiceProvider();

    services.AddMassTransit(configure =>
    {
        configure.SetKebabCaseEndpointNameFormatter();
        configure.AddConsumers(assembly);

        configure.UsingRabbitMq((context, configurator) =>
        {
            var rabbitSettings = serviceProvider.GetService<IOptions<RabbitSettings>>().Value;
            var host = new Uri("rabbitmq://" + rabbitSettings.EventBusConnection);

            configurator.Host(host, h =>
            {
                h.Username(rabbitSettings.EventBusUserName);
                h.Password(rabbitSettings.EventBusPassword);
            });

            var types = AppDomain.CurrentDomain.GetAssemblies().SelectMany(x => x.GetTypes())
                .Where(x => x.BaseType == typeof(IntegrationEvent));

            foreach (var type in types)
            {
                var consumers = AppDomain.CurrentDomain.GetAssemblies().SelectMany(x => x.GetTypes())
                    .Where(x => x.IsAssignableTo(typeof(IConsumer<>).MakeGenericType(type))).ToList();

                if (consumers.Any())
                {
                    // rabbitSettings.QueueName => service-b
                    configurator.ReceiveEndpoint(rabbitSettings.QueueName, e =>
                        {
                            e.UseConsumeFilter(typeof(InboxFilter<>), context);
                            foreach (var consumer in consumers)
                            {
                                configurator.ConfigureEndpoints(context, x => x.Exclude(consumer));

                                var methodInfo = typeof(DependencyInjectionReceiveEndpointExtensions)
                                    .GetMethods()
                                    .Where(x => x.GetParameters()
                                        .Any(p => p.ParameterType == typeof(IServiceProvider)))
                                    .FirstOrDefault(x => x.Name == "Consumer" && x.IsGenericMethod);

                                var generic = methodInfo?.MakeGenericMethod(consumer);
                                generic?.Invoke(e, new object[] { e, context, null });
                            }
                        });
                }
            }
        });
    });

    return services;
}

IntegrationEvent是所有集成事件的基类,已从拓扑中排除。

问题分析

原配置的核心问题:

  • 手动通过反射调用Consumer方法添加消费者的逻辑冗余且错误,导致MassTransit无法正确关联消费者与事件
  • 循环内重复调用configurator.ConfigureEndpoints并排除消费者,破坏了MassTransit的自动配置流程
  • 未利用框架内置的消费者配置扩展,导致队列无法正常创建

可行解决方案

优化后的配置简化了消费者动态发现逻辑,直接使用MassTransit内置的ConfigureConsumers方法完成消费者与接收端点的绑定,同时增加了熔断、重试等可靠性配置:

public static IServiceCollection AddCustomMassTransit(this IServiceCollection services)
{
    var serviceProvider = services.BuildServiceProvider();

    services.AddMassTransit(configure =>
    {
        configure.SetKebabCaseEndpointNameFormatter();

        var allTypes = AppDomain.CurrentDomain.GetAssemblies().SelectMany(x => x.GetTypes()).ToArray();
        var eventTypes = allTypes.Where(x => x.BaseType == typeof(IntegrationEvent)).ToArray();
        Type[] consumerTypes = allTypes.Where(x => eventTypes.Any(et => x.IsAssignableTo(typeof(IConsumer<>).MakeGenericType(et)))).ToArray();

        configure.AddConsumers(consumerTypes);

        configure.UsingRabbitMq((context, configurator) =>
        {
            var rabbitSettings = serviceProvider.GetService<IOptions<RabbitSettings>>().Value;
            var host = new Uri("rabbitmq://" + rabbitSettings.EventBusConnection);

            configurator.Host(host, h =>
            {
                h.Username(rabbitSettings.EventBusUserName);
                h.Password(rabbitSettings.EventBusPassword);
            });

            configurator.UseCircuitBreaker(cb =>
            {
                cb.TrackingPeriod = TimeSpan.FromMinutes(1);
                cb.TripThreshold = 15;
                cb.ActiveThreshold = 10;
                cb.ResetInterval = TimeSpan.FromMinutes(5);
            });

            configurator.UseMessageRetry(r =>
            {
                r.Ignore(typeof(ArgumentException),
                         typeof(ArgumentNullException),
                         typeof(ArgumentOutOfRangeException),
                         typeof(IndexOutOfRangeException),
                         typeof(DivideByZeroException),
                         typeof(InvalidCastException));

                r.Intervals(new[] { 1, 2, 4, 8, 16 }.Select(t => TimeSpan.FromSeconds(t)).ToArray());
            });

            if (consumerTypes.Length > 0)
            {
                configurator.ReceiveEndpoint(rabbitSettings.QueueName, e =>
                {
                    e.UseConsumeFilter(typeof(InboxFilter<>), context);
                    e.ConfigureConsumers(context);
                });
            }
        });
    });

    return services;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 13:36:37