缓存发布订阅模式实现中泛型委托Action<T>转Action<string>失败问题排查
解决缓存Pub/Sub模式中不同类型Action回调的兼容问题
你的问题根源在于当前设计把所有类型的订阅者都绑定到同一个Action<string>事件上,导致发布的JSON只能被完全匹配类型的订阅者正确反序列化,其他类型的订阅者在反序列化时会抛出异常(比如把字符串"Value1"转int肯定失败),最终回调无法执行。而且这种设计违背了Pub/Sub模式中"同类型消息匹配订阅者"的常规逻辑——int类型的订阅者本就不该接收string类型的消息。
下面给你两种针对性的解决方案,根据你的实际场景选择:
方案一:内存内Pub/Sub(无需序列化,性能最优)
如果你的Pub/Sub只在内存内运行,完全不需要序列化,直接按类型分组管理订阅者即可,这样既能保证类型安全,又能避免序列化开销:
public class CacheNotification { // 按类型存储对应回调列表,键是消息类型,值是该类型的所有回调 private readonly Dictionary<Type, List<Delegate>> _subscribers = new(); public Task Publish<T>(T value) { var messageType = typeof(T); // 只通知订阅了当前消息类型的回调 if (_subscribers.TryGetValue(messageType, out var callbacks)) { // 用ToList避免遍历过程中集合被修改(比如取消订阅) foreach (var callback in callbacks.ToList()) { if (callback is Action<T> typedCallback) { typedCallback(value); } } } return Task.CompletedTask; } public Task Subscribe<T>(Action<T> callback) { var messageType = typeof(T); if (!_subscribers.ContainsKey(messageType)) { _subscribers[messageType] = new List<Delegate>(); } _subscribers[messageType].Add(callback); return Task.CompletedTask; } // 可选:添加取消订阅方法,完善生命周期管理 public Task Unsubscribe<T>(Action<T> callback) { var messageType = typeof(T); if (_subscribers.TryGetValue(messageType, out var callbacks)) { callbacks.Remove(callback); if (callbacks.Count == 0) { _subscribers.Remove(messageType); } } return Task.CompletedTask; } }
测试代码调整
你的原测试逻辑有问题:发布string类型的消息,int类型的订阅者本就不该被触发。如果要测试不同类型的订阅,应该分别发布对应类型的消息:
int sub1Counter = 0; int sub2Counter = 0; Action<string> sub1Callback = msg => sub1Counter++; Action<int> sub2Callback = num => sub2Counter++; var cache = new CacheNotification(); await cache.Subscribe(sub1Callback); await cache.Subscribe(sub2Callback); // 发布string消息,触发sub1 await cache.Publish("Value1"); Assert.Equal(1, sub1Counter); Assert.Equal(0, sub2Counter); // 这里应该是0,因为没发布int消息 // 发布int消息,触发sub2 await cache.Publish(123); Assert.Equal(1, sub1Counter); Assert.Equal(1, sub2Counter);
方案二:需要序列化的场景(比如跨进程/跨机器通信)
如果你的Pub/Sub需要把消息序列化后传输(比如发送到消息队列),那需要在发布时携带类型信息,订阅时先校验类型再反序列化:
1. 定义带类型信息的消息结构
public class PublishedMessage { // 用AssemblyQualifiedName可以准确找到类型 public string TypeFullName { get; set; } public string JsonContent { get; set; } }
2. 修改CacheNotification类
public class CacheNotification { private readonly Dictionary<Type, List<Delegate>> _subscribers = new(); // 对外暴露的序列化事件,用于跨进程传输 public Action<PublishedMessage> OnSerializedMessagePublished; public Task Publish<T>(T value) { var messageType = typeof(T); // 1. 通知内存内的强类型订阅者(无需序列化) if (_subscribers.TryGetValue(messageType, out var callbacks)) { foreach (var callback in callbacks.ToList()) { if (callback is Action<T> typedCallback) { typedCallback(value); } } } // 2. 序列化消息并触发对外事件(用于跨进程传输) var json = JsonSerializer.Serialize(value); var serializedMsg = new PublishedMessage { TypeFullName = messageType.AssemblyQualifiedName, JsonContent = json }; OnSerializedMessagePublished?.Invoke(serializedMsg); return Task.CompletedTask; } public Task Subscribe<T>(Action<T> callback) { var messageType = typeof(T); if (!_subscribers.ContainsKey(messageType)) { _subscribers[messageType] = new List<Delegate>(); } _subscribers[messageType].Add(callback); return Task.CompletedTask; } // 处理从外部接收的序列化消息(比如从队列读取后) public Task ProcessSerializedMessage(PublishedMessage message) { var messageType = Type.GetType(message.TypeFullName); if (messageType == null || !_subscribers.TryGetValue(messageType, out var callbacks)) { return Task.CompletedTask; } // 反序列化并调用对应类型的回调 var value = JsonSerializer.Deserialize(message.JsonContent, messageType); foreach (var callback in callbacks.ToList()) { callback.DynamicInvoke(value); } return Task.CompletedTask; } }
可选:支持类型转换的订阅(特殊场景)
如果你确实需要让int订阅者接收string消息并自动转换,可以添加一个支持类型转换的订阅方法:
public Task Subscribe<TInput, TOutput>(Action<TOutput> callback, Func<TInput, TOutput> converter) { var inputType = typeof(TInput); if (!_subscribers.ContainsKey(inputType)) { _subscribers[inputType] = new List<Delegate>(); } // 包装回调:先转换类型再执行 Action<TInput> wrappedCallback = input => callback(converter(input)); _subscribers[inputType].Add(wrappedCallback); return Task.CompletedTask; } // 使用示例:订阅string消息,转换为int后执行回调 await cache.Subscribe<string, int>( num => sub2Counter++, s => int.TryParse(s, out var n) ? n : 0 // 处理转换失败的情况 );
这样发布"123"时,int订阅者就能正常触发了。
内容的提问来源于stack exchange,提问作者Haha
相关产品推荐
相关产品推荐

