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

重复订阅同一Redis频道时,如何避免接收重复通知?

重复订阅Redis频道导致接收重复消息的解决方法

问题场景

重复调用订阅方法订阅同一Redis频道时,订阅者会收到重复消息,相关代码如下:

订阅方法

public void SubscribeToUserNotifications(string userId, Action<string?> onMessage)
{
    _subscriber.Subscribe
        (
        new RedisChannel($"orders:notifications:user:{userId}",
            RedisChannel.PatternMode.Literal), (channel, message) =>
            {
                onMessage(message);
            });
}

消息发布方法

public async Task AddOrderForUserAsync(string userId, OrderInfo order)
{
    var db = _redis.GetDatabase();

    string orderKey = $"orders:user:{userId}";

    string orderData = JsonConvert.SerializeObject(order);

    await db.HashSetAsync(orderKey, order.Id.ToString(), orderData);

    _subscriber.Publish(
        new RedisChannel($"orders:notifications:user:{userId}", RedisChannel.PatternMode.Literal),
        orderData
    );
}

核心问题:每次调用订阅方法都会创建同一频道的新订阅,导致同一条消息被多次接收。Redis没有直接避免重复订阅的内置机制,需实现"即使重复订阅,每条消息仅接收一次"的效果。

解决方案

Redis本身没有内置的重复订阅去重机制,必须在应用层处理,常见实现方式如下:

1. 维护已订阅频道的本地缓存

在应用内用字典记录每个用户ID对应的订阅状态,避免重复发起Redis订阅:

private readonly Dictionary<string, Action<string?>> _userSubscriptions = new Dictionary<string, Action<string?>>();
private readonly object _subscriptionLock = new object();

public void SubscribeToUserNotifications(string userId, Action<string?> onMessage)
{
    lock (_subscriptionLock)
    {
        // 检查是否已订阅该用户频道
        if (_userSubscriptions.ContainsKey(userId))
        {
            // 可选:替换旧回调或直接返回
            _userSubscriptions[userId] = onMessage;
            return;
        }

        // 未订阅则执行Redis订阅
        _subscriber.Subscribe(
            new RedisChannel($"orders:notifications:user:{userId}", RedisChannel.PatternMode.Literal),
            (channel, message) => onMessage(message)
        );
        _userSubscriptions.Add(userId, onMessage);
    }
}
  • 加锁保证线程安全,避免并发订阅的竞态问题
  • 若需支持取消订阅,需补充移除逻辑并调用_subscriber.Unsubscribe

2. 基于消息ID兜底去重

如果无法完全避免重复订阅,可在发布消息时携带唯一标识(如订单ID),订阅端维护已接收消息ID列表过滤重复:

// 订阅端的消息处理逻辑
private readonly HashSet<string> _receivedOrderIds = new HashSet<string>();
private readonly object _idLock = new object();

public void HandleOrderNotification(string message)
{
    var order = JsonConvert.DeserializeObject<OrderInfo>(message);
    if (order == null) return;

    lock (_idLock)
    {
        string orderId = order.Id.ToString();
        if (_receivedOrderIds.Contains(orderId))
        {
            // 已处理过该消息,直接返回
            return;
        }
        _receivedOrderIds.Add(orderId);
    }

    // 执行实际业务逻辑
    ProcessOrderNotification(order);

    // 可选:定期清理过期ID,避免内存溢出
    CleanupOldOrderIds();
}
  • 适合无法完全控制订阅次数的场景,作为兜底方案
  • 分布式环境下可改用Redis过期键存储已接收ID,实现全局去重

3. 单例订阅者统一管理

将订阅逻辑封装为单例服务,确保每个频道仅被订阅一次,所有需要接收消息的组件通过回调注册到该服务:

public class OrderNotificationService
{
    private readonly IRedisSubscriber _subscriber;
    private readonly Dictionary<string, List<Action<string?>>> _channelCallbacks = new Dictionary<string, List<Action<string?>>>();
    private readonly object _lock = new object();

    public OrderNotificationService(IRedisSubscriber subscriber)
    {
        _subscriber = subscriber;
    }

    public void Subscribe(string userId, Action<string?> callback)
    {
        string channel = $"orders:notifications:user:{userId}";
        lock (_lock)
        {
            if (!_channelCallbacks.ContainsKey(channel))
            {
                // 首次订阅该频道
                _subscriber.Subscribe(new RedisChannel(channel, RedisChannel.PatternMode.Literal), (ch, msg) =>
                {
                    foreach (var cb in _channelCallbacks[ch])
                    {
                        cb(msg);
                    }
                });
                _channelCallbacks.Add(channel, new List<Action<string?>>());
            }
            // 将回调添加到频道的回调列表
            _channelCallbacks[channel].Add(callback);
        }
    }

    // 取消订阅方法,移除对应回调并清理空频道的Redis订阅
    public void Unsubscribe(string userId, Action<string?> callback)
    {
        string channel = $"orders:notifications:user:{userId}";
        lock (_lock)
        {
            if (_channelCallbacks.TryGetValue(channel, out var callbacks))
            {
                callbacks.Remove(callback);
                if (callbacks.Count == 0)
                {
                    _subscriber.Unsubscribe(new RedisChannel(channel, RedisChannel.PatternMode.Literal));
                    _channelCallbacks.Remove(channel);
                }
            }
        }
    }
}
  • 统一管理所有订阅,从根源避免重复订阅同一频道
  • 适合多组件需要接收同一频道消息的场景

内容的提问来源于stack exchange,提问作者Gleb Stepanov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 01:55:58