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

RSocket-JS客户端刷新后无法接收Spring RSocket服务端Flux推送消息如何解决

问题根因分析
  • FluxSink引用未正确更新:你全局复用了taskUpdateNoticeConsumer实例,第一次请求流触发Flux.create()时,consumer会持有本次请求对应的FluxSink引用,后续新连接发起请求时,新的Flux.create传入的新sink没有被consumer替换,导致所有事件都发给已经断开的旧连接sink,新连接对应的sink根本没有收到事件。你观测到的next调用大概率还是在调用旧sink的next方法,只是旧连接断开后没有抛出异常而已。
  • 热流操作符使用不当:你用的publish().autoConnect()在第一个订阅者连接后就会永久启动上游,就算所有订阅者断开也不会重置上游状态,新订阅者只能收到订阅之后的事件;如果用的share()(即publish().refCount(1)),订阅者数量清零时会取消上游,但如果你的taskUpdateNoticeConsumer没有处理重订阅逻辑,第二次订阅时也无法正常接收事件。
  • RSocket连接生命周期绑定错误:服务端的流没有和当前RSocket连接的生命周期绑定,连接断开时没有正确清理对应的sink引用。
解决步骤

1. 重构事件通知逻辑,支持多连接订阅

不要让consumer只持有单个sink引用,改成持有线程安全的sink集合,每个新连接请求流时把新sink加入集合,连接断开时自动移除对应sink,事件下发时遍历所有存活的sink发送:

// 全局定义线程安全的sink集合
private final Set<FluxSink<TaskBrief>> taskSinks = new CopyOnWriteArraySet<>();

// 事件发布逻辑统一修改为遍历所有存活sink
public void publishTaskEvent(TaskBrief event) {
    taskSinks.forEach(sink -> {
        if (!sink.isCancelled()) {
            sink.next(event);
        }
    });
}

2. 调整服务端控制器实现

不需要每次调用都加publish/autoConnect/share操作符,直接返回和当前连接生命周期绑定的独立Flux即可:

@MessageMapping("request.stream")
public Flux<String> requestStream(@Payload String userInfo) throws JsonProcessingException {
    return Flux.<TaskBrief>create(sink -> {
        // 新连接sink加入集合
        taskSinks.add(sink);
        // 连接断开/取消订阅时自动移除sink
        sink.onDispose(() -> taskSinks.remove(sink));
    })
    .map(taskEvent -> {
        // 保留原有业务转换逻辑
        return ...;
    });
}

3. 修复客户端流订阅销毁逻辑

rsocket-js 0.0.23版本需要显式取消流订阅再关闭连接,避免服务端sink残留:

useEffect(() => {
    let streamSub = null;
    async function initRsocketConnection() {
        await client.connect().subscribe({        
            onComplete: socket => {
                const stream = socket.requestStream({
                    data: "jsclient",
                    metadata: String.fromCharCode("request.stream".length) + "request.stream"
                })
                streamSub = stream.subscribe({
                    onNext: (data: Payload<unknown, unknown>) => {
                        // 原有消息处理逻辑
                    },
                    onError: err => console.error(err),
                    onSubscribe: sub => {
                        sub.request(2147483647);
                    }
                });
            }
        });
    }

    initRsocketConnection();
    return () => {
        // 先取消流订阅再关闭客户端
        if (streamSub) streamSub.cancel();
        client.close();
    }
}, [showUpdates])

4. 可选配置调整

如果需要给新连接的订阅者重放历史消息,可以在服务端Flux.create后添加replay操作符,比如.replay(10)就是给新订阅者重放最近10条历史消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 21:39:03