迁移至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
相关产品推荐
相关产品推荐

