【发布时间】: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->{}, 0).doSomethingMore()之类的吗?我对 reactor 和 reactor-netty 很陌生。因此,更大的解释将非常有帮助。谢谢你。
标签: spring-webflux project-reactor reactor-netty