Spring Reactive Redis Pub/Sub如何取消订阅(断开频道连接)?
Spring Reactive Redis 取消订阅方案
一、针对ReactiveRedisTemplate.listenToChannel()的处理方式
你当前代码中,listenToChannel()返回的是Flux流,调用subscribe()会得到一个Disposable对象——这是控制Redis订阅生命周期的核心。socket通道关闭后不会自动触发Redis订阅取消,必须主动调用Disposable.dispose()来终止订阅并释放资源。
修改步骤:
- 维护订阅与socket通道的关联:用线程安全的容器存储每个
ChannelHandlerContext对应的Disposable,方便通道关闭时快速定位订阅。 - 监听socket通道关闭事件:给通道添加关闭监听器,通道关闭时立即调用对应
Disposable的dispose()方法,取消Redis订阅。 - 消息发送失败时同步取消订阅:在消息发送失败的回调逻辑中,除了关闭通道,还要主动取消Redis订阅,避免无效订阅残留。
修改后的代码示例:
@Service class RedisMessageSubscriber( private val redisTemplate: ReactiveRedisTemplate<String, Any>, private val reactiveRedisMessageListenerContainer: ReactiveRedisMessageListenerContainer ) { private val log = logger() // 线程安全Map,存储通道与订阅Disposable的关联关系 private val channelSubscriptionMap = ConcurrentHashMap<ChannelHandlerContext, Disposable>() fun subscribe(channel: String, ctx: ChannelHandlerContext) { val subscription = redisTemplate.listenToChannel(channel) .map { MapperUtils.readJsonValueOrThrow(it.message.toString(), ChatDto::class) } .subscribe( { chatDto -> ctx.channel().writeAndFlush(CustomResponse(ResponseType.RECEIVE_CHAT, chatDto)) .addListener { future -> if (!future.isSuccess) { ctx.channel().close() log.warn("Failed to send the message to the client. roomId:${channel}, message:${future.cause()}, socketOpen: ${ctx.channel().isOpen}") // 发送失败时主动取消订阅 channelSubscriptionMap.remove(ctx)?.dispose() } else { log.info("write channel topic: $channel") } } }, // 处理Redis订阅流异常,避免资源泄漏 { error -> log.error("Redis subscription error for channel: $channel", error) channelSubscriptionMap.remove(ctx)?.dispose() ctx.channel().close() } ) // 关联通道与订阅 channelSubscriptionMap[ctx] = subscription log.info("isDispose: ${subscription.isDisposed}") // 监听通道关闭事件,触发订阅取消 ctx.channel().closeFuture().addListener { channelSubscriptionMap.remove(ctx)?.dispose() log.info("Socket channel closed, cancelled Redis subscription for channel: $channel") } } }
二、针对ReactiveRedisMessageListenerContainer的处理方式
如果使用ReactiveRedisMessageListenerContainer进行订阅,订阅操作会返回Subscription对象,取消订阅时直接调用subscription.cancel()即可。
示例代码:
fun subscribeWithContainer(channel: String, ctx: ChannelHandlerContext) { val channelTopic = ChannelTopic(channel) val messageFlux = reactiveRedisMessageListenerContainer.receive(channelTopic) .map { it.message } val subscription = messageFlux.subscribe( { message -> // 业务消息处理逻辑 }, { error -> // 异常处理逻辑 } ) // 绑定通道关闭事件与订阅取消 ctx.channel().closeFuture().addListener { subscription.cancel() log.info("Cancelled subscription via container for channel: $channel") } }
关键注意点
- Reactor流生命周期:Reactive Redis的订阅本质是Reactor的
Flux流,仅当流完成、出错,或者主动调用Disposable.dispose()/Subscription.cancel()时,才会终止Redis订阅并释放底层连接资源。 - 线程安全:Netty通道事件与Redis订阅可能在不同线程执行,存储关联关系的容器必须保证线程安全(如
ConcurrentHashMap)。
内容的提问来源于stack exchange,提问作者tray
相关产品推荐
相关产品推荐

