如何实现强类型消息Broker:反序列化消息并动态调用运行时处理器
嘿,你这个强类型消息Broker的需求我太熟悉了——之前做类似的事件驱动系统时踩过不少坑,刚好能给你一套可行的方案,完美解决你提到的三个问题。
核心思路
我们要做的是启动时一次性扫描所有DI组件的处理器方法,建立「消息类型→强类型委托」的映射,然后在分发时直接根据消息类型反序列化并调用对应的强类型处理器,全程保持类型安全,避免dynamic的无提示、反射的性能损耗和基类转换的瓶颈。
分步实现
1. 定义基础结构
先把核心的特性、消息基类(可选但推荐)定义好:
// 标记处理器方法的特性 [AttributeUsage(AttributeTargets.Method, AllowMultiple = false)] public class HandlesAttribute : Attribute { public Type MessageType { get; } public HandlesAttribute(Type messageType) { MessageType = messageType; } } // 消息基类(可选,用来统一消息标识比如MessageId) public abstract class Message { public Guid MessageId { get; set; } = Guid.NewGuid(); } // 示例消息 public class ExampleMessage : Message { public string Content { get; set; } = string.Empty; }
2. 编写带处理器的组件
组件完全是黑盒,通过DI注入依赖,用[Handles]标记处理器方法(支持私有/公有方法):
public class ExampleComponent { private readonly ILogger<ExampleComponent> _logger; // DI注入依赖,完全符合黑盒要求 public ExampleComponent(ILogger<ExampleComponent> logger) { _logger = logger; } [Handles(typeof(ExampleMessage))] private void HandleExampleMessage(ExampleMessage message) { _logger.LogInformation($"处理消息[{message.MessageId}]: {message.Content}"); // 这里写你的业务逻辑 } // 也支持异步处理器 [Handles(typeof(AnotherMessage))] private async Task HandleAnotherMessageAsync(AnotherMessage message) { await Task.Delay(100); _logger.LogInformation("异步处理完成"); } }
3. 实现消息Broker核心
这部分是关键:启动时扫描DI组件,建立类型映射,分发时直接调用强类型委托:
public interface IMessageBroker { Task DispatchAsync(string serializedMessage, CancellationToken cancellationToken = default); } public class MessageBroker : IMessageBroker { // 存储「消息类型→强类型处理器委托」的映射(支持一个消息多个处理器) private readonly Dictionary<Type, List<Delegate>> _messageHandlers = new(); public MessageBroker(IServiceProvider serviceProvider) { // 从DI容器中获取所有注册的服务,扫描处理器方法 var allServices = serviceProvider.GetServices<object>(); foreach (var service in allServices) { RegisterHandlersFromService(service); } } private void RegisterHandlersFromService(object service) { var serviceType = service.GetType(); // 筛选所有带[Handles]特性的方法(支持私有/公有) var handlerMethods = serviceType.GetMethods( BindingFlags.Instance | BindingFlags.Public | BindingFlags.NonPublic) .Where(m => m.GetCustomAttribute<HandlesAttribute>() != null); foreach (var method in handlerMethods) { var handlesAttr = method.GetCustomAttribute<HandlesAttribute>()!; var messageType = handlesAttr.MessageType; // 验证处理器签名:必须只有一个对应消息类型的参数 var parameters = method.GetParameters(); if (parameters.Length != 1 || parameters[0].ParameterType != messageType) { throw new InvalidOperationException( $"处理器方法 {serviceType.Name}.{method.Name} 签名无效:必须接受单个 {messageType.Name} 参数"); } // 创建强类型委托(同步用Action<T>,异步用Func<T, Task>) var delegateType = method.ReturnType == typeof(Task) ? typeof(Func<,>).MakeGenericType(messageType, typeof(Task)) : typeof(Action<>).MakeGenericType(messageType); var handler = Delegate.CreateDelegate(delegateType, service, method); // 注册到映射表 if (!_messageHandlers.ContainsKey(messageType)) { _messageHandlers[messageType] = new List<Delegate>(); } _messageHandlers[messageType].Add(handler); } } public async Task DispatchAsync(string serializedMessage, CancellationToken cancellationToken = default) { // 第一步:从序列化消息中解析出消息类型 // 假设序列化消息包含"MessageType"字段,用来指定消息的完整类型名 var jsonDoc = JsonDocument.Parse(serializedMessage); var messageTypeFullName = jsonDoc.RootElement.GetProperty("MessageType").GetString(); var messageType = Type.GetType(messageTypeFullName); if (messageType == null) { throw new InvalidOperationException($"找不到消息类型:{messageTypeFullName}"); } // 第二步:反序列化为强类型消息 var message = JsonSerializer.Deserialize( serializedMessage, messageType, new JsonSerializerOptions { PropertyNameCaseInsensitive = true }); if (message == null) { throw new InvalidOperationException("反序列化消息失败"); } // 第三步:找到对应的处理器并调用 if (!_messageHandlers.TryGetValue(messageType, out var handlers)) { throw new InvalidOperationException($"没有找到 {messageType.Name} 的处理器"); } foreach (var handler in handlers) { if (handler is Func<object, Task> asyncHandler) { await asyncHandler(message).WaitAsync(cancellationToken); } else if (handler is Action<object> syncHandler) { syncHandler(message); } } } }
4. DI注册与使用
在启动时注册你的组件和消息Broker:
// Program.cs var builder = WebApplication.CreateBuilder(args); // 注册你的黑盒组件(支持Scoped/Transient/Singleton) builder.Services.AddScoped<ExampleComponent>(); builder.Services.AddScoped<AnotherComponent>(); // 注册消息Broker(Singleton或Scoped都可以,根据你的需求) builder.Services.AddSingleton<IMessageBroker, MessageBroker>(); var app = builder.Build(); // 使用示例:比如在API接口中接收序列化消息并分发 app.MapPost("/messages", async (string serializedMessage, IMessageBroker broker) => { await broker.DispatchAsync(serializedMessage); return Results.Ok(); }); app.Run();
为什么这个方案能解决你的问题?
- 强类型安全:处理器方法直接接收强类型消息,编译时就能检查参数错误,IDE也有智能提示,完全避免了
dynamic的问题。 - 反射仅启动时执行:所有处理器的扫描和委托创建只在应用启动时做一次,运行时分发完全是直接调用委托,性能优异,代码也整洁。
- 避免基类转换瓶颈:用
Delegate存储强类型委托,不需要把消息转成Message基类再强制转换回来,彻底解决了第三个方案的类型转换问题。
额外优化建议
- 消息类型解析优化:可以用
AssemblyQualifiedName来存储消息类型,或者维护一个类型映射配置表,避免Type.GetType找不到跨程序集类型的问题。 - 错误处理:在
DispatchAsync中添加try-catch,捕获处理器抛出的异常,做日志记录、重试或者死信队列处理。 - 批量分发:如果需要处理批量消息,可以扩展方法支持批量反序列化和调用。
内容的提问来源于stack exchange,提问作者Matthew Goulart
相关产品推荐
相关产品推荐

