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

基于Reactor-Netty的WebSocket:如何等待变量变更后再发送内容

在Reactor-Netty中实现WebSocket时等待变量变更后发送内容的解决方案

当你在Reactor-Netty里实现WebSocket,需要等待某个变量(比如消息集合)变更后再发送下一条内容时,我们可以通过消息缓存集合+读取刷新方法的组合来实现,既能等待新内容,又不会立即关闭WebSocket连接。

核心思路

我们用一个集合来临时存储待发送的消息,当需要发送时,通过专门的方法取出当前所有消息并清空集合,这样就能保证每次发送的都是最新变更后的内容;同时利用Reactor的响应式特性,让WebSocket连接保持活跃,等待下一次的消息变更。

代码实现

private Set<String> message = new HashSet<>();

// 添加待发送的消息到缓存集合
private void writeMessage(String message) {
    this.message.add(message);
}

// 读取并清空当前缓存的消息,返回待发送的内容数组
private String[] readFlushMessage() {
    String[] _message = (String[]) this.message.toArray();
    this.message = new HashSet<>();
    return _message;
}

// WebSocket发布器实现,处理连接并发送消息
private Publisher<Void> websocketPublisherA(HttpServerRequest request, HttpServerResponse response, WebSocketServerHandle handleObject) {
    return response
            .header("content-type", "text/plain")
            .sendWebsocket((in, out) -> out.options(NettyPipeline.SendOptions::flushOnEach)
                    .sendString(Flux.just(readFlushMessage())));
}

关键代码解释

  • message集合:作为消息的临时缓存,当外部有新消息需要发送时,调用writeMessage()将消息加入集合,完成“变量变更”的触发。
  • readFlushMessage():这个方法是核心,它会一次性取出当前所有缓存的消息,然后重置集合为空,确保下一次读取的是全新的消息内容。
  • sendWebsocket()中的逻辑:通过out.options(NettyPipeline.SendOptions::flushOnEach)设置每条消息都立即刷新,Flux.just(readFlushMessage())将读取到的消息转为响应式流发送,同时因为连接没有主动关闭,会保持打开状态等待下一次的消息变更触发。

内容的提问来源于stack exchange,提问作者Henry

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:20:28