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

缓存发布订阅模式实现中泛型委托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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 08:32:30