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

Spring Boot RSocket如何捕获并处理Cancel帧信号?

如何捕获并处理RSocket的Cancel信号

首先得明确:你用@MessageMapping("cancel")或@ConnectMapping("cancel")抓不到信号是正常的——这两个注解是用来处理应用层自定义消息的,而RSocket协议本身的Cancel帧是属于Reactive流层面的取消信号,得用Reactive的回调机制来捕获。

接下来给你具体的实现方案:

1. 在业务流中直接绑定取消回调

RSocket的Cancel帧会触发对应Reactive流(比如Flux或Mono)的cancel()调用,所以你只需要在@MessageMapping方法返回的流上添加doOnCancel()操作符,就能在取消发生时执行你的清理逻辑。

举个请求流(Request-Stream)的例子:

@MessageMapping("data.stream")
public Flux<Data> streamData(RequestParams params) {
    // 假设这是你的业务数据流
    return dataService.continuousDataStream(params)
        // 绑定取消回调,在这里处理其他订阅的取消、资源清理
        .doOnCancel(() -> {
            log.info("RSocket Cancel信号已触发,开始清理资源");
            // 取消服务器上的其他订阅注册
            subscriptionRegistry.cancelAllForClient(params.getClientId());
            // 释放相关资源
            resourceManager.release(params.getResourceKey());
        });
}

不管是客户端主动发送Cancel帧,还是WebSocket连接断开触发的Cancel,这个回调都会被执行——因为WebSocket关闭会触发RSocket连接终止,进而触发所有活跃流的取消。

2. 全局拦截所有RSocket请求的Cancel信号

如果需要对所有RSocket请求统一处理Cancel逻辑,可以自定义一个RSocket拦截器,通过Spring Boot的ServerRSocketFactoryCustomizer来注册:

@Bean
public ServerRSocketFactoryCustomizer globalCancelHandler() {
    return factory -> factory.addInterceptor(new RSocketInterceptor() {
        @Override
        public RSocket intercept(RSocket targetRSocket, Metadata metadata) {
            return new RSocket() {
                // 针对请求-响应模式
                @Override
                public Mono<Payload> requestResponse(Payload payload) {
                    return targetRSocket.requestResponse(payload)
                        .doOnCancel(() -> log.info("Request-Response 请求已被取消"));
                }

                // 针对请求-流模式
                @Override
                public Flux<Payload> requestStream(Payload payload) {
                    return targetRSocket.requestStream(payload)
                        .doOnCancel(() -> {
                            log.info("Request-Stream 请求已被取消");
                            // 在这里添加全局的取消处理逻辑
                        });
                }

                // 其他交互模式(Fire-and-Forget、Channel)同理,按需重写
                @Override
                public Mono<Void> fireAndForget(Payload payload) {
                    return targetRSocket.fireAndForget(payload)
                        .doOnCancel(() -> log.info("Fire-and-Forget 请求已被取消"));
                }

                // 若需要监听连接级别的断开,还可以重写dispose方法
                @Override
                public void dispose() {
                    targetRSocket.dispose();
                    log.info("RSocket连接已断开,执行全局清理");
                }
            };
        }
    });
}

为什么之前的方法没用?

再啰嗦一句:@MessageMapping("cancel")是用来接收客户端主动发送的自定义消息(比如客户端专门发一个名为"cancel"的应用层消息),而你日志里看到的cancel()是RSocket框架处理协议帧时触发的Reactive流取消,完全不是一回事,所以根本不会路由到这个映射端点。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 20:27:51