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

Micronaut中Flux超时无法中断WebSocket连接尝试问题

问题:Micronaut 3.9.0 WebSocket连接超时控制失效

在Micronaut 3.9.0中开发WebSocket连接时,尝试通过Flux的timeout操作符设置连接超时,但超时触发后连接尝试并未中断,最终仍成功建立了连接。

示例测试代码:

package com.example;

import io.micronaut.context.BeanContext;
import io.micronaut.test.extensions.junit5.annotation.MicronautTest;
import io.micronaut.websocket.WebSocketClient;
import io.micronaut.websocket.WebSocketSession;
import io.micronaut.websocket.annotation.ClientWebSocket;
import io.micronaut.websocket.annotation.OnClose;
import io.micronaut.websocket.annotation.OnMessage;
import io.micronaut.websocket.annotation.OnOpen;
import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Assertions;

import jakarta.inject.Inject;
import org.reactivestreams.Publisher;
import reactor.core.publisher.Flux;

import java.time.Duration;
import java.util.concurrent.atomic.AtomicBoolean;

@Slf4j
@MicronautTest
class WstimeoutTest {
    static AtomicBoolean isOpened = new AtomicBoolean(false);

    @Inject
    BeanContext beanContext;

    @ClientWebSocket
    static abstract class TestWebSocketClient implements AutoCloseable {
        @OnMessage
        void onMessage(String message) {
            log.info("[onMessage] Got message: {}", message);
        }

        @OnOpen
        void onOpen(WebSocketSession session) {
            log.info("[onOpen] WS session id: {}", session.getId());
            isOpened.set(true);
        }

        @OnClose
        void onClose(WebSocketSession session) {
            log.info("[onClose] WS session id: {}", session.getId());
        }
    }

    @Test
    void testWebSocketTimeout() {
        var uri = "ws://localhost";
        var timeout = Duration.ofMillis(40); // 设置预期的连接超时时间

        var webSocketClient = beanContext.getBean(WebSocketClient.class);

        Publisher<TestWebSocketClient> client = webSocketClient.connect(TestWebSocketClient.class,  uri);
        Flux.from(client)
                .timeout(timeout)
                .doOnError(throwable -> log.info("Expected error: {}", throwable.getMessage()))
                .subscribe();


        Awaitility.await().atLeast(Duration.ofMillis(250)).untilAsserted(
                () -> Assertions.assertFalse(isOpened.get(), "WebSocket should not be opened")
        );
    }
}

测试输出:

2023:11:21T13:24:56.140 [parallel-1] INFO  com.example.WstimeoutTest[][] - Expected error: Did not observe any item or terminal signal within 40ms in 'switchMapNoPrefetch' (and no fallback has been configured)
2023:11:21T13:24:56.178 [default-nioEventLoopGroup-1-2] INFO com.example.WstimeoutTest[][] - [onOpen] WS session id: AOKppbvUUl5zc3/x1N0lBiXW4WU=

WebSocket should not be opened
Expected :false
Actual   :true
解决方案

问题核心是timeout操作符仅终止了上层Flux流,但没有中断底层的WebSocket连接尝试。要彻底终止连接,需要在超时发生时主动清理底层资源,或从配置层面控制连接超时。

方案1:超时错误中主动关闭WebSocket客户端

修改订阅逻辑,在超时触发时关闭客户端,终止所有未完成的连接尝试:

@Test
void testWebSocketTimeout() {
    var uri = "ws://localhost";
    var timeout = Duration.ofMillis(40);

    var webSocketClient = beanContext.getBean(WebSocketClient.class);

    Flux.from(webSocketClient.connect(TestWebSocketClient.class, uri))
            .timeout(timeout)
            .doOnError(throwable -> {
                log.info("Expected error: {}", throwable.getMessage());
                // 关闭客户端,终止底层连接尝试
                webSocketClient.close();
            })
            .doOnNext(TestWebSocketClient::close) // 若连接成功也主动关闭
            .subscribe();

    Awaitility.await().atLeast(Duration.ofMillis(250)).untilAsserted(
            () -> Assertions.assertFalse(isOpened.get(), "WebSocket should not be opened")
    );
}

方案2:通过配置文件设置底层超时

在application.yml中配置WebSocket客户端的连接和读取超时,从底层限制连接尝试时长:

micronaut:
  http:
    client:
      connect-timeout: 40ms
      read-timeout: 40ms
    websocket:
      client:
        ping-interval: 0 # 禁用ping避免干扰超时逻辑

该配置会作用于所有WebSocket客户端连接,无需在代码中手动处理超时。

方案3:用Reactor资源管理自动清理

使用Flux.using管理WebSocket客户端生命周期,确保流终止(包括超时)时自动关闭资源:

@Test
void testWebSocketTimeout() {
    var uri = "ws://localhost";
    var timeout = Duration.ofMillis(40);

    Flux.using(
            // 资源提供者:创建WebSocketClient
            () -> beanContext.getBean(WebSocketClient.class),
            // 业务逻辑:发起连接并设置超时
            client -> Flux.from(client.connect(TestWebSocketClient.class, uri))
                    .timeout(timeout),
            // 资源清理:自动关闭客户端
            WebSocketClient::close
    )
    .doOnError(throwable -> log.info("Expected error: {}", throwable.getMessage()))
    .subscribe();

    Awaitility.await().atLeast(Duration.ofMillis(250)).untilAsserted(
            () -> Assertions.assertFalse(isOpened.get(), "WebSocket should not be opened")
    );
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 16:42:35