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

