如何确保Reactor-Netty WebSocketSession前次send完成再发消息并规避慢客户端问题
问题解决方案
一、确保会话前一次send操作完成后再发送新消息
你当前的代码里,同一会话的多次send操作是并发执行的,没有保证顺序和前一次操作的完成性;而且复用同一个stringFlux给多个会话发送,会导致每个会话订阅时重新遍历消息集合(虽然功能正常,但不够高效)。
解决思路是给每个会话维护一个发送序列的Mono引用,每次新的发送请求都通过then串联到上一次发送的Mono之后,确保串行执行:
- 调整会话存储结构,新增一个Map记录每个会话的当前发送任务:
// 存储活跃会话 private final Map<String, WebSocketSession> sessions = new ConcurrentHashMap<>(); // 存储每个会话的待完成发送任务,初始值为Mono.empty() private final Map<String, Mono<Void>> sessionSendTasks = new ConcurrentHashMap<>();
- 发送消息时的核心逻辑:
Flux<String> stringFlux = Flux.fromIterable(messages); for (Map.Entry<String, WebSocketSession> entry : sessions.entrySet()) { String sessionId = entry.getKey(); WebSocketSession session = entry.getValue(); if (!session.isOpen()) { System.out.println("session is closed.. skipping.. " + sessionId); sessions.remove(sessionId); sessionSendTasks.remove(sessionId); continue; } // 创建本次发送任务:确保整个消息流发送完成后才标记结束 Mono<Void> currentSend = session.send(stringFlux.map(session::textMessage)) .then(); // 原子更新发送任务:将本次任务串联到上一次任务之后,保证串行执行 sessionSendTasks.compute(sessionId, (id, previousTask) -> { return previousTask == null ? currentSend : previousTask.then(currentSend); }).subscribe( () -> {}, error -> { System.err.println("send failed for session " + sessionId + ": " + error.getMessage()); sessions.remove(sessionId); sessionSendTasks.remove(sessionId); } ); }
这样每个会话的所有发送操作都会按顺序执行,前一次send完成后才会启动下一次。
二、避免向慢客户端写入导致内存开销
Reactor Netty和Spring WebFlux本身提供了多种机制处理慢客户端问题:
1. 利用内置背压机制
WebSocketSession.send()返回的Mono<Void>天然遵守背压规则:当客户端读取速度跟不上服务器发送速度时,底层Netty通道的写缓冲区会被填满,此时send操作会自动暂停发送,直到客户端读取部分数据、缓冲区有空闲空间,不会持续占用内存。
你可以通过配置调整缓冲区的水位线,优化内存控制:
// 配置Reactor Netty服务器时设置写缓冲区参数 HttpServer.create() .tcpConfiguration(tcp -> tcp .bootstrap(serverBootstrap -> serverBootstrap .option(ChannelOption.WRITE_BUFFER_HIGH_WATER_MARK, 32 * 1024) // 高水位:32KB .option(ChannelOption.WRITE_BUFFER_LOW_WATER_MARK, 16 * 1024) // 低水位:16KB ) );
当缓冲区达到高水位时,Netty会停止接收新的写请求,直到缓冲区降到低水位以下,避免内存持续增长。
2. 添加发送超时
如果客户端长时间不读取数据导致发送操作阻塞,可以给send添加超时逻辑,超时后主动关闭会话:
Mono<Void> currentSend = session.send(stringFlux.map(session::textMessage)) .then() .timeout(Duration.ofSeconds(30)) // 设置30秒超时 .doOnError(TimeoutException.class, error -> { System.err.println("send timeout for session " + sessionId + ", closing session"); session.close().subscribe(); });
3. 定期清理无响应会话
通过定时任务检查会话的发送任务状态,若任务长时间处于未完成状态,判定为慢客户端并主动清理:
// 定时任务,每分钟执行一次 ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(() -> { sessionSendTasks.forEach((sessionId, task) -> { if (!task.isDisposed() && !task.isSuccess()) { WebSocketSession session = sessions.get(sessionId); if (session != null && session.isOpen()) { System.err.println("slow client detected, closing session " + sessionId); session.close().subscribe(); sessions.remove(sessionId); sessionSendTasks.remove(sessionId); } } }); }, 1, 1, TimeUnit.MINUTES);
内容的提问来源于stack exchange,提问作者Amudhan
相关产品推荐
相关产品推荐

