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

如何基于有限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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 18:02:50