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

Azure ServiceBus是否有NServiceBus.IHandleMessages等效功能?需多队列吗?

解决Azure.Messaging.ServiceBus多消息类型重复代码问题

不用为每种消息类型单独创建队列或处理器,Azure.Messaging.ServiceBus作为底层SDK确实没有NServiceBus那种开箱即用的消息路由,但你可以自己实现类似的集中式消息分发机制,核心思路是共用队列+类型路由:

核心方案

1. 共用队列/主题订阅

所有消息类型都发送到同一个队列(或主题的同一个订阅),无需拆分队列。

2. 给消息标记类型标识

发送消息时,在消息的ApplicationProperties中添加类型标识,比如:

var message = new ServiceBusMessage(JsonSerializer.Serialize(orderCreatedEvent))
{
    ApplicationProperties = { ["MessageType"] = typeof(OrderCreatedEvent).FullName }
};

或者序列化消息时带上类型信息(比如用Newtonsoft.Json的TypeNameHandling.All),反序列化时直接识别类型。

3. 实现通用处理器+消息路由

  • 定义统一的消息处理接口,模仿NServiceBus的IHandleMessage:
    public interface IMessageHandler<T>
    {
        Task HandleAsync(T message, CancellationToken cancellationToken);
    }
    
  • 为每种消息类型实现这个接口:
    public class OrderCreatedHandler : IMessageHandler<OrderCreatedEvent>
    {
        public async Task HandleAsync(OrderCreatedEvent message, CancellationToken cancellationToken)
        {
            // 业务处理逻辑
        }
    }
    
  • 注册所有处理器到依赖注入容器(以ASP.NET Core为例):
    // 扫描程序集内所有IMessageHandler<T>实现并注册
    var handlerTypes = AppDomain.CurrentDomain.GetAssemblies()
        .SelectMany(a => a.GetTypes())
        .Where(t => t.GetInterfaces().Any(i => i.IsGenericType && i.GetGenericTypeDefinition() == typeof(IMessageHandler<>)));
    
    foreach (var type in handlerTypes)
    {
        var interfaceType = type.GetInterfaces().First(i => i.IsGenericType && i.GetGenericTypeDefinition() == typeof(IMessageHandler<>));
        services.AddScoped(interfaceType, type);
    }
    
  • 编写通用ServiceBus处理器,接收所有消息并路由到对应处理器:
    public class GenericMessageProcessor
    {
        private readonly IServiceProvider _serviceProvider;
    
        public GenericMessageProcessor(IServiceProvider serviceProvider)
        {
            _serviceProvider = serviceProvider;
        }
    
        public async Task ProcessMessageAsync(ProcessMessageEventArgs args)
        {
            // 从消息属性获取类型标识
            if (!args.Message.ApplicationProperties.TryGetValue("MessageType", out var typeNameObj) || 
                typeNameObj is not string typeName)
            {
                await args.DeadLetterMessageAsync(args.Message, "MissingMessageType", "未找到消息类型标识");
                return;
            }
    
            var messageType = Type.GetType(typeName);
            if (messageType == null)
            {
                await args.DeadLetterMessageAsync(args.Message, "UnknownMessageType", $"未找到类型 {typeName}");
                return;
            }
    
            // 反序列化消息体
            var messageBody = await args.Message.Body.ToObjectFromJsonAsync(messageType);
    
            // 获取对应处理器实例
            var handlerInterface = typeof(IMessageHandler<>).MakeGenericType(messageType);
            using var scope = _serviceProvider.CreateScope();
            var handler = scope.ServiceProvider.GetService(handlerInterface);
            if (handler == null)
            {
                await args.DeadLetterMessageAsync(args.Message, "MissingHandler", $"未找到 {typeName} 的处理器");
                return;
            }
    
            // 调用处理方法
            var handleMethod = handlerInterface.GetMethod(nameof(IMessageHandler<object>.HandleAsync));
            await (Task)handleMethod.Invoke(handler, new[] { messageBody, args.CancellationToken });
    
            // 标记消息处理完成
            await args.CompleteMessageAsync(args.Message);
        }
    }
    
  • 注册通用处理器到ServiceBus客户端:
    var processor = client.CreateProcessor("shared-queue", new ServiceBusProcessorOptions());
    var genericProcessor = serviceProvider.GetRequiredService<GenericMessageProcessor>();
    processor.ProcessMessageAsync += genericProcessor.ProcessMessageAsync;
    processor.ProcessErrorAsync += (args) =>
    {
        // 错误处理逻辑
        return Task.CompletedTask;
    };
    await processor.StartProcessingAsync();
    

要不要继续用NServiceBus?

如果你的项目需要以下特性,直接用NServiceBus更省心:

  • 开箱即用的消息路由、重试、死信处理
  • 分布式事务、Saga(长流程)支持
  • 成熟的监控、生态集成能力
  • 团队已熟悉NServiceBus开发模式

如果只是需要轻量消息收发,不想依赖商业框架,自己实现上述路由机制完全可以解决重复代码问题,且灵活性更高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 10:05:11