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

Spring Reactive Redis Pub/Sub如何取消订阅(断开频道连接)?

Spring Reactive Redis 取消订阅方案

一、针对ReactiveRedisTemplate.listenToChannel()的处理方式

你当前代码中,listenToChannel()返回的是Flux流,调用subscribe()会得到一个Disposable对象——这是控制Redis订阅生命周期的核心。socket通道关闭后不会自动触发Redis订阅取消,必须主动调用Disposable.dispose()来终止订阅并释放资源。

修改步骤:

  1. 维护订阅与socket通道的关联:用线程安全的容器存储每个ChannelHandlerContext对应的Disposable,方便通道关闭时快速定位订阅。
  2. 监听socket通道关闭事件:给通道添加关闭监听器,通道关闭时立即调用对应Disposable的dispose()方法,取消Redis订阅。
  3. 消息发送失败时同步取消订阅:在消息发送失败的回调逻辑中,除了关闭通道,还要主动取消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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 06:05:08