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
相关产品推荐
相关产品推荐

