重复订阅同一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
相关产品推荐
相关产品推荐

