如何基于有限Flux构建无限Flux?WebSocket断连无需重绑订阅
问题分析与解决方案
咱们先拆解下你的核心痛点:WebSocket断连后,原来的Flux会完成,导致所有绑定的执行器都得重新订阅——这确实是Reactor构建长连接系统时的常见问题,你的临时Holder方案虽然能跑,但确实不够优雅,而且容易引入线程安全、订阅时机错位之类的潜在问题。下面给你梳理几个更符合Reactor最佳实践的方案:
1. 用Flux.retryWhen实现自动重连的无限Flux
这是处理长连接断连恢复的标准姿势,完全不需要手动管理FluxSink或者EmitterProcessor。核心思路是:把单次WebSocket连接的逻辑包装在Flux.defer里,然后用retryWhen处理断连后的重试,让整个流变成“断开就自动重建”的无限流,订阅者只需要订阅一次,断连后会自动续上。
重构你的WebSocket实现如下:
public Flux<HaEventResponse> startListening(Flux<HaEventRequest> eventRequests) { return Flux.defer(() -> { // 每次重连时创建全新的登录+订阅流 Mono<String> login = Mono.just(loginPayload()); Flux<String> subscribe = eventRequests .doOnNext(req -> log.info("Registering event: {}", req.getEventType())) .map(Json::write); Flux<String> input = login.concatWith(subscribe); return client.execute(URI.create(wsUrl), session -> session.send(input.map(session::textMessage)) .thenMany(session.receive() .map(WebSocketMessage::getPayloadAsText) .map(message -> Json.read(message, HaEventResponse.class)) .filter(event -> event.getEvent() != null)) .doOnSubscribe(s -> log.info("WebSocket connection established")) .doOnError(e -> log.error("WebSocket connection failed", e)) .doOnComplete(() -> log.info("WebSocket connection closed")) ); }) // 指数退避重试:最多重试5次,每次间隔1秒,加入随机抖动避免雪崩 .retryWhen(Retry.backoff(5, Duration.ofSeconds(1)) .jitter(0.5) .doBeforeRetry(signal -> log.info("Retrying WebSocket, attempt {}", signal.totalRetries() + 1))); } // 订阅示例:只需要订阅一次,断连后自动恢复 public void bindLightTrigger() { startListening(eventRequestsFlux) .filter(event -> "button_pressed".equals(event.getEvent().getEventType())) .subscribe(event -> { // 执行开灯逻辑 log.info("Turning on light!"); }); }
这个方案的优势:
- 订阅者无需感知断连,一次订阅永久生效
- 完全符合Reactor声明式编程风格,避免手动状态管理
- 内置重试策略(指数退避、抖动),比手动重试更健壮
2. 关于Processor的使用场景
如果你的需求不是单纯的重连,而是需要缓存事件、支持新订阅者获取历史数据,那可以结合ReplayProcessor,但它只是补充,不是替代重连逻辑的核心:
private final ReplayProcessor<HaEventResponse> eventCache = ReplayProcessor.create(10); // 缓存最近10个事件 public void startListening(Flux<HaEventRequest> eventRequests) { Flux.defer(() -> { // 同上面的连接逻辑 }) .retryWhen(Retry.backoff(5, Duration.ofSeconds(1))) .subscribe(eventCache::onNext, eventCache::onError); } public Flux<HaEventResponse> streamEvents(String eventType) { return eventCache.filter(event -> event.getEvent().getEventType().equals(eventType)); }
但注意:Processor需要额外处理线程安全(比如用toSerialized()),而且对于单纯的重连场景,直接用retryWhen更简洁。
3. 解决“必须提前知晓所有事件类型”的问题
原来的实现需要在startListening时传入所有事件请求,限制了动态订阅的灵活性。可以用Subject来实现动态添加订阅:
// 线程安全的Subject,用于接收动态的事件订阅请求 private final Subject<HaEventRequest> eventRequestSubject = UnicastProcessor.create().toSerialized(); // 对外提供动态订阅接口 public void subscribeToEventType(HaEventRequest eventRequest) { eventRequestSubject.onNext(eventRequest); // 优化:如果当前WebSocket连接存在,立即发送订阅请求(需要额外维护连接状态) } public Flux<HaEventResponse> startListening() { return Flux.defer(() -> { // 重连时自动订阅所有当前已注册的事件 Flux<HaEventRequest> eventRequests = eventRequestSubject.replay().autoConnect(); Mono<String> login = Mono.just(loginPayload()); Flux<String> subscribe = eventRequests .doOnNext(req -> log.info("Registering event: {}", req.getEventType())) .map(Json::write); Flux<String> input = login.concatWith(subscribe); return client.execute(URI.create(wsUrl), session -> session.send(input.map(session::textMessage)) .thenMany(session.receive() .map(WebSocketMessage::getPayloadAsText) .map(message -> Json.read(message, HaEventResponse.class)) .filter(event -> event.getEvent() != null)) ); }) .retryWhen(Retry.backoff(5, Duration.ofSeconds(1))); }
这样你可以在运行时随时添加新的事件订阅,不需要提前知晓所有类型。
总结
- 放弃你的Holder临时方案:
Flux.retryWhen+Flux.defer是Reactor处理长连接自动重连的标准最佳实践,代码更简洁、健壮。 - Processor不是必须的:只有在需要事件缓存、历史回放等特殊场景下才用,普通重连场景用操作符足够。
- 动态订阅需求用
Subject解决:避免提前绑定所有事件类型的限制。
内容的提问来源于stack exchange,提问作者p.streef
相关产品推荐
相关产品推荐

