【问题标题】:How to make sure the previous session.send() is complete before sending more data in reactor-netty WebSocketSession如何在 reactor-netty WebSocketSession 中发送更多数据之前确保之前的 session.send() 已完成
【发布时间】:2022-11-22 20:51:38
【问题描述】:

我正在使用 spring-webflux 和 reactor-netty 创建一个 WebSocket 服务器,并且在每个连接上,我都会得到一个 WebSocketSession(带有 ReactorNettyWebSocketSession 实例)。假设我将所有会话存储在一个映射中,并想根据一些业务逻辑向这个 WebSocketSessions 发送一些消息,我可以做类似的事情

            Flux<String> stringFlux = Flux.fromIterable(messages);
            for (WebSocketSession session : sessions.values()) {
                if (session.isOpen()) {
                    session
                            .send(stringFlux.map(session::textMessage))
                            .subscribe();
                } else {
                    System.out.println("session is closed.. skipping.. " + session.getId());
                    sessions.remove(session.getId());
                }
            }

现在,当我向会话发送消息时,有没有办法确保会话之前的发送已完成?如果客户端速度非常慢和/或不从服务器读取,如果服务器在客户端不读取或读取速度非常慢时继续写入套接字,则可能会为服务器产生内存开销。

如果客户端速度很慢,如何获得回调或某种机制来阻止写入 WebSocketSession/socket?

【问题讨论】:

  • 您可以使用 concatMap 并将预取设置为 0 以实现此目的。不会像想象的那样使用您的代码。考虑将会话值也设为 Flux。
  • @Khepu,对不起,我不明白。你是说Flux.fromIterable(sessions.values()).concatMap(f-&gt;{}, 0).doSomethingMore()之类的吗?我对 reactor 和 reactor-netty 很陌生。因此,更大的解释将非常有帮助。谢谢你。

标签: spring-webflux project-reactor reactor-netty


【解决方案1】:

这个答案背后有很多假设,所以我们可能需要在 cmets 中讨论细节。这个想法是以下代码为每个会话一个接一个地发送消息,但同时为所有会话发送消息。您可以根据您的需要进行调整,将concatMap 改为flatMap 将启用并发发送,而将flatMap 改为concatMap 将确保所有消息都必须发送到一个会话,然后才能继续下一个会话一。

Flux.fromIterable(sessions.values())
    .filter(session -> session.isOpen()) // keep open sessions
    .doOnDsicard(Session.class, session -> {
        System.out.println("session is closed.. skipping.. " + session.getId());
        sessions.remove(session.getId());
    })
    // flatMap here will allow the messages to be sent to each socket concurrently
    .flatMap(session -> Flux.fromIterable(messages) 
        .map(session::textMessage)
        // concatMap here will force messages to be sent one by one (per session)
        .concatMap(message -> Mono
            .fromCallable(() -> session
                .send(Mono.just(message))), 0)) 

【讨论】:

    猜你喜欢
    • 2019-08-10
    • 1970-01-01
    • 2019-06-13
    • 1970-01-01
    • 1970-01-01
    • 2018-01-23
    • 2013-11-21
    • 2020-02-16
    • 1970-01-01
    相关资源
    最近更新 更多