不能直接用 SUBSCRIBE+回调处理高频消息,因其独占连接且回调阻塞I/O线程、缺乏背压机制,易致消息积压和内存溢出;应改用 Channel + IAsyncEnumerable 实现可控消费与流控。

为什么不能直接用 SUBSCRIBE + 回调处理高频消息
StackExchange.Redis 的 SUBSCRIBE 方法底层会独占一个连接,并通过回调(Action<redischannel redisvalue></redischannel>)推送消息。这种模式在低频场景下没问题,但一旦每秒消息量超过几百条,就会暴露两个硬伤:一是回调执行在 Redis 客户端的 I/O 线程上,阻塞该线程会导致后续消息积压甚至连接超时;二是无法天然支持背压(backpressure),订阅者消费慢时,消息会在内存中无节制堆积,最终触发 OutOfMemoryException。
IAsyncEnumerable<T> 是解耦与流控的关键
真正适合高频 Pub/Sub 的方式,是把订阅通道“转成”异步流,让业务层用 await foreach 按需拉取,而不是被动接收。核心在于:不依赖客户端回调,而是用独立的后台任务持续读取 ChannelReader<RedisValue>,再封装为 IAsyncEnumerable<T>。
关键步骤包括:
- 用
ConnectionMultiplexer.GetSubscriber()获取ISubscriber实例,但**不调用SUBSCRIBE** - 改用
Channel<RedisValue>作为中间缓冲,启动一个长期运行的Task调用subscriber.Subscribe(channelName).OnMessage(...),并在回调里向Channel.Writer.WriteAsync(...)写入 - 对外暴露
Channel.Reader.ReadAllAsync()返回的IAsyncEnumerable<RedisValue>,业务代码可自由控制消费节奏 - 务必设置
Channel.CreateBounded<RedisValue>(capacity: 1024),避免无限缓存
如何避免 Channel 关闭后消息丢失
Redis Pub/Sub 是“发即忘”机制:如果订阅者断连或未及时读取,消息就彻底丢失。这和 Stream 不同,没有重放能力。所以必须在应用层做兜底,否则高频下丢消息是常态。
Redis 缓存和数据结构管理技能。通过自然语言操作 Redis,支持 String、Hash、List、Set、ZSet、Stream 等数据结构操作。当用户提到 Redis、缓存、消息队列、会话存储时使用此技能。
常见应对策略有:
- 在
OnMessage回调中,对写入Channel.Writer的操作加try/catch,捕获ChannelClosedException后主动重连并重新SUBSCRIBE - 不要依赖
Channel.Reader.Completion判断订阅是否终止——它只反映本地 Channel 状态,不反映 Redis 连接是否存活 - 定期用
PING命令探测 Redis 连接健康度,发现异常时主动重建ISubscriber和Channel - 若业务要求不丢消息,必须放弃纯 Pub/Sub,改用
Stream+XREADGROUP,哪怕牺牲一点实时性
实际消费时最容易忽略的反模式
很多人以为只要用了 await foreach 就算“异步流”,但下面这些写法会让整个流退化为同步阻塞:
- 在
await foreach循环体内直接调用同步 IO 方法(如File.WriteAllText、HttpClient.Send),阻塞当前线程 - 没用
ConfigureAwait(false),导致 await 后回调被强行调度回 UI/ASP.NET 上下文,引发死锁或性能抖动 - 在循环内创建新
HttpClient实例,触发 DNS 查询和连接池竞争 - 把消息解析逻辑(如
JsonSerializer.Deserialize<Event>(msg))放在循环内,而没预热JsonSerializerOptions或复用Utf8JsonReader
高频 Pub/Sub 的瓶颈往往不在 Redis,而在你消费消息的那一小段代码里。别让反模式吃掉所有吞吐优势。

















