如何在take(1)等条件触发时完全关闭Reactor WebSocket连接
问题原因
代码无限挂起的核心原因有两个:
- 入站接收流单独调用了
subscribe(),它是独立的订阅生命周期,执行完take(1)结束后不会通知WebSocket连接需要关闭 - 出站逻辑调用了
neverComplete(),强制告知框架保持连接活性,永远不会主动触发连接关闭逻辑,导致连接一直存活。
解决方法
将入站、出站的处理逻辑组合成同一个响应式流,当入站满足终止条件后整个流完成,Reactor Netty会自动断开WebSocket连接,程序正常退出。
修改后的可正常退出的代码如下:
import io.netty.buffer.Unpooled; import io.netty.util.CharsetUtil; import reactor.core.publisher.Flux; import reactor.netty.http.client.HttpClient; public class Application { public static void main(String[] args) { HttpClient client = HttpClient.create(); client.websocket() .uri("wss://echo.websocket.org") .handle((inbound, outbound) -> { final byte[] msgBytes = "hello".getBytes(CharsetUtil.ISO_8859_1); // 出站逻辑:发送两条消息 var sendFlux = outbound.send(Flux.just(Unpooled.wrappedBuffer(msgBytes), Unpooled.wrappedBuffer(msgBytes))); // 入站逻辑:接收1条消息后终止流 var receiveFlux = inbound.receive() .asString() .take(1) .doOnNext(System.out::println); // 组合两个流:先完成消息发送,再等待入站处理完成,流结束后自动关闭连接 return sendFlux.then(receiveFlux); }) .blockLast(); } }
自定义终止条件扩展
如果需要用自定义逻辑判断何时关闭连接,只需要修改入站流的终止操作符即可:
- 用
takeWhile(消息 -> 保留连接的条件):当条件不满足时终止入站流 - 用
takeUntil(消息 -> 关闭连接的条件):当条件满足时立即终止入站流
如果你的场景需要边发边收、不需要等发送完成再处理入站,也可以用Flux.merge(sendFlux, receiveFlux).then()来组合流,只要任意一个流完成,整体就会终止并关闭连接。
内容的提问来源于stack exchange,提问作者Alec
相关产品推荐
相关产品推荐

