Spring Reactive WebSocket:多Hot Stream消息推送异常排查
问题分析与解决方案
我帮你梳理下代码里的核心问题,以及对应的修复方案——你的问题主要出在流的生命周期管理和系统消息的触发逻辑上,导致热流的新事件无法持续推送到客户端。
核心问题拆解
系统消息的触发逻辑错误
你现在的relevantSystemMessages是通过帖子流的元素去触发过滤系统消息,这意味着只有当有新帖子推送时,才会去消费系统消息的当前事件。如果帖子流停止(比如历史消息推送完、没有触发加载更多),系统消息的订阅也会跟着停滞,完全没法实时推送新的系统消息。多流推送的生命周期冲突
你用.and()组合了三个独立的session.send(),但每个send()对应的Flux如果是有限流(比如recentPosts拉完历史就结束),对应的Mono<Void>会提前完成,可能影响其他热流的订阅生命周期。流的重复订阅隐患
recentPosts被同时用在session.send()和Flux.merge()中,如果它是冷流,会导致重复拉取数据;即使是热流,分散的订阅也会增加生命周期管理的复杂度。
修复后的代码
@Bean public WebSocketHandler webSocketHandler(MessageProducer<Message> messageProducer) { return session -> { // 1. 合并所有帖子流:先推历史帖子,再推加载更多的帖子 Flux<Message> recentPosts = messageProducer.streamRecentPosts(); Flux<Message> onLoadMorePosts = session.receive() .map(WebSocketMessage::getPayloadAsText) .map(this::toInputMessage) .filter(LoadMoreInputMessage.class::isInstance) .map(LoadMoreInputMessage.class::cast) .map(LoadMoreInputMessage::getLoaded) .flatMap(messageProducer::streamRecentPosts); Flux<Message> allPosts = Flux.concat(recentPosts, onLoadMorePosts); // 2. 维护已推送的帖子ID集合,用于过滤系统消息 // 用ReplayProcessor缓存所有已推送的帖子ID,确保新系统消息能匹配全部历史帖子 ReplayProcessor<String> postIdsCache = ReplayProcessor.create(); allPosts.map(Message::getId).subscribe(postIdsCache); // 用scan维护动态更新的帖子ID集合 Flux<Set<String>> currentPostIds = postIdsCache.scan(HashSet::new, (idSet, postId) -> { idSet.add(postId); return idSet; }); // 3. 持续监听系统消息,过滤出需要推送的内容 Flux<Message> relevantSystemMessages = messageProducer.streamSystemMessages() .withLatestFrom(currentPostIds, (systemMsg, ids) -> { // 全局消息(无关联帖子ID)或关联已推送帖子的消息才推送 if (systemMsg.getAffectedPostId() == null || ids.contains(systemMsg.getAffectedPostId())) { return systemMsg; } return null; }) .filter(Objects::nonNull); // 4. 合并所有消息流,统一发送给客户端 Flux<Message> allMessages = Flux.merge(allPosts, relevantSystemMessages); return session.send(allMessages.map(this::stringify).map(session::textMessage)); }; }
关键修复点说明
- 合并帖子流:用
Flux.concat确保先推送历史帖子,再处理加载更多的请求,同时避免重复订阅的问题。 - 缓存已推送帖子ID:通过
ReplayProcessor和scan维护动态更新的帖子ID集合,让系统消息能持续匹配所有已推送的帖子,不再依赖帖子流的活动。 - 统一消息发送:把所有要推送的消息合并成一个流,用单次
session.send()处理,确保所有热流的新事件都能被持续订阅和推送。 - 系统消息持续监听:用
withLatestFrom让每个新系统消息都和当前最新的帖子ID集合做匹配,保证系统消息能实时推送。
另外,还要确保你的MessageProducer中的各个streamXXX方法正确返回热流,比如用DirectProcessor或ReplayProcessor实现,并在有新消息时调用sink.next()推送事件,示例如下:
@Component public class MessageProducer { private final DirectProcessor<Message> systemMsgProcessor = DirectProcessor.create(); private final FluxSink<Message> systemMsgSink = systemMsgProcessor.sink(); public Flux<Message> streamSystemMessages() { return systemMsgProcessor.share(); // share保证多订阅者共享同一个热流 } // 外部调用这个方法推送新系统消息 public void pushSystemMessage(Message msg) { systemMsgSink.next(msg); } // 帖子相关的Processor实现同理... }
内容的提问来源于stack exchange,提问作者Dmytro
相关产品推荐
相关产品推荐

