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

迁移至Azure.Messaging.ServiceBus:如何保留反射解析泛型消息策略

问题

从Microsoft.ServiceBus的旧版client.OnMessageAsync()迁移至Azure.Messaging.ServiceBus的client.CreateProcessor()时,原消息处理器通过反射解析泛型消息,如何在新版SDK中实现相同逻辑?

旧版实现代码:

public class AzureBusReceiver<T>
{
        public async Task OnReceived(BrokeredMessage brokeredMessage, T processor)
                {
                    var messageType = Type.GetType(brokeredMessage.Properties["messageType"].ToString());
                    var method = typeof(BrokeredMessage).GetMethod("GetBody", new Type[] { });
                    if (method != null)
                    {
                        var generic = method.MakeGenericMethod(messageType);
                        var messageBody = generic.Invoke(brokeredMessage, null);
                        var args = new[] { messageBody };
                        try
                        {
                            await new DynamicProcessor<T>().Run(processor, messageType, args);
                        }
                        catch
                        {
                            // ignored
                        }
                    }
                }
    }
解决方案

在Azure.Messaging.ServiceBus中,核心思路是替换消息类型、调整消息体解析方式,保留原有的反射调用逻辑:

关键变化

  • 旧版BrokeredMessage替换为新版ServiceBusReceivedMessage
  • 不再通过GetBody<T>方法获取消息体,改为读取ServiceBusReceivedMessage.Body的二进制内容并反序列化
  • 反射调用处理器的逻辑可直接复用

新版实现代码

using Azure.Messaging.ServiceBus;
using System.IO;
using System.Runtime.Serialization.Formatters.Binary;

public class AzureBusReceiver<T>
{
    public async Task OnReceived(ServiceBusReceivedMessage receivedMessage, T processor)
    {
        // 从消息属性中提取目标类型
        if (!receivedMessage.ApplicationProperties.TryGetValue("messageType", out var typeValue) 
            || typeValue is not string typeName)
        {
            return;
        }

        var targetType = Type.GetType(typeName);
        if (targetType == null)
        {
            return;
        }

        // 反序列化消息体(沿用旧版BinaryFormatter,可替换为JSON等其他序列化方式)
        object messageBody;
        using (var stream = new MemoryStream(receivedMessage.Body.ToArray()))
        {
            var formatter = new BinaryFormatter();
            messageBody = formatter.Deserialize(stream);
        }

        // 调用原有动态处理器逻辑
        try
        {
            await new DynamicProcessor<T>().Run(processor, targetType, new[] { messageBody });
        }
        catch
        {
            // 保留原异常处理逻辑,建议根据实际需求补充日志或重试
        }
    }
}

处理器注册示例

将上述方法绑定到ServiceBusProcessor的消息处理事件:

var serviceBusClient = new ServiceBusClient("your-connection-string");
var processor = serviceBusClient.CreateProcessor("your-queue-name");

// 注册消息处理回调
processor.ProcessMessageAsync += async args =>
{
    var receiver = new AzureBusReceiver<YourProcessorImpl>();
    await receiver.OnReceived(args.Message, new YourProcessorImpl());
    
    // 标记消息已处理完成
    await args.CompleteMessageAsync(args.Message);
};

// 启动处理器
await processor.StartProcessingAsync();

注意事项

  • 如果旧版使用的不是BinaryFormatter(比如JSON序列化),需要将消息体反序列化逻辑替换为对应的序列化器(如System.Text.Json.JsonSerializer)。
  • 建议补充异常处理逻辑,不要直接忽略异常,可添加日志记录或死信处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 04:14:55