如何通过Micronaut WebSocket端点发送无界Flux更新并避免异常
问题:WebSocket客户端断开连接时Flux发送操作抛出异常
我有一个按固定间隔发射无界数量值的Flux,需要一个WebSocket端点在客户端连接时发送该Flux的值。当前实现如下:
@ServerWebSocket("/updates") public class UpdateController { private Flux<Update> updates; // ... 为简洁起见省略部分代码 @OnOpen public Flux<Update> onOpen(WebSocketSession session) { return updates.flatMap(session::send); } @OnMessage public void onMessage(String content) { // do nothing } @OnClose public void onClose(WebSocketSession session) { // do nothing } }
该实现可正常工作,但客户端断开连接时会抛出异常。我理解这是因为updates Flux仍在发射值,且会调用session::send方法。请问该如何调整代码结构来避免这个异常?
解决方案
核心问题是客户端断开后,WebSocketSession已不可用,但无界的updates Flux仍在持续发射数据,调用session.send()会触发异常。需要让发送操作随会话关闭自动终止,并优雅处理异常。
调整后的代码实现
@ServerWebSocket("/updates") public class UpdateController { private Flux<Update> updates; // ... 省略初始化代码 @OnOpen public Flux<Void> onOpen(WebSocketSession session) { return updates // 将Update转换为WebSocket文本消息(可根据实际需求改为二进制消息) .map(update -> session.textMessage(update.toString())) // 发送消息 .flatMap(session::send) // 监听会话关闭信号,一旦会话关闭立即终止Flux订阅 .takeUntilOther(session.closeStatus()) // 捕获发送过程中的异常,避免异常扩散 .onErrorResume(e -> Mono.empty()); } @OnMessage public void onMessage(String content) { // 无需处理消息可保留空实现 } @OnClose public void onClose(WebSocketSession session) { // 会话关闭逻辑已通过takeUntilOther处理,无需额外操作 } }
关键修改说明
- 修正返回类型:
session.send()返回Mono<Void>,flatMap后整体Flux类型改为Flux<Void>,符合语义。 - 绑定生命周期到会话:
takeUntilOther(session.closeStatus())会监听会话的关闭信号,一旦客户端断开,自动终止updates的订阅,停止后续发送操作。 - 异常优雅处理:
onErrorResume捕获发送阶段的异常(比如会话突然断开),避免异常向上传播导致报错。 - 显式消息转换:明确将
Update对象转换为WebSocketMessage,避免隐式转换可能带来的问题。
内容的提问来源于stack exchange,提问作者AutomatedMess
相关产品推荐
相关产品推荐

