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

如何确保Reactor-Netty WebSocketSession前次send完成再发消息并规避慢客户端问题

问题解决方案

一、确保会话前一次send操作完成后再发送新消息

你当前的代码里,同一会话的多次send操作是并发执行的,没有保证顺序和前一次操作的完成性;而且复用同一个stringFlux给多个会话发送,会导致每个会话订阅时重新遍历消息集合(虽然功能正常,但不够高效)。

解决思路是给每个会话维护一个发送序列的Mono引用,每次新的发送请求都通过then串联到上一次发送的Mono之后,确保串行执行:

  1. 调整会话存储结构,新增一个Map记录每个会话的当前发送任务:
// 存储活跃会话
private final Map<String, WebSocketSession> sessions = new ConcurrentHashMap<>();
// 存储每个会话的待完成发送任务,初始值为Mono.empty()
private final Map<String, Mono<Void>> sessionSendTasks = new ConcurrentHashMap<>();
  1. 发送消息时的核心逻辑:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 09:10:31