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

Micronaut ServerWebSocket异步消息发送问题求助

理解背压错误的含义

RxJava的**背压(Backpressure)**是用来解决生产者发送数据速度超过消费者处理速度的机制。MissingBackpressureException在这里的意思是:你调用session.send()返回的Publisher在发射消息时,没有收到来自订阅者的请求,或者消费者处理速度跟不上生产者的发送节奏,导致消息无法被正常处理。

在Micronaut的@OnOpen、@OnMessage等注解方法中,框架会自动帮你订阅返回的Publisher,并处理背压逻辑,所以你直接返回session.send(...)就能正常工作。但当你在这些方法外部手动调用session.send()时,需要自己处理订阅和背压,否则就会出现问题。

异步发送消息的正确实现方式

下面给你两种贴合Micronaut+RxJava生态的可靠实现方案:

方案1:用RxJava的Flowable.interval替代ScheduledExecutorService

这种方式更符合RxJava的编程模型,能自动处理调度和背压:

@CompileStatic
@ServerWebSocket("/ws")
class MyWebSocket {
    @OnOpen
    public Publisher<String> onOpen(WebSocketSession session) {
        println("opening connection")
        startPing(session)
        return session.send("Hello")
    }

    @OnMessage
    public Publisher<String> onMessage(String message, WebSocketSession session) {
        println("message received: $message")
        return session.send("Thanks for the message: $message")
    }

    @OnClose
    public void onClose(WebSocketSession session) {
        println("Closing connection")
        // 会话关闭时无需发送消息,客户端已感知连接断开
    }

    private void startPing(WebSocketSession session) {
        // 每隔1秒发送一次Ping消息
        Flowable.interval(0, 1, TimeUnit.SECONDS)
            // 切换到IO线程池执行发送,避免阻塞Netty事件循环
            .subscribeOn(Schedulers.io())
            // 发送消息并添加背压缓存策略
            .flatMap { tick ->
                session.send("Ping... $tick")
                    // 缓存待发送消息,直到WebSocket能处理
                    .onBackpressureBuffer()
            }
            // 订阅并处理异常
            .subscribe(
                { /* 发送成功的回调,可留空 */ },
                { error -> 
                    println("Ping发送失败: ${error.message}")
                    error.printStackTrace()
                }
            )
    }
}

方案2:若坚持用ScheduledExecutorService,需正确订阅Publisher

如果你更习惯Java调度器,记得手动订阅session.send()返回的Publisher并处理背压:

@CompileStatic
@ServerWebSocket("/ws")
class MyWebSocket {
    @OnOpen
    public Publisher<String> onOpen(WebSocketSession session) {
        println("opening connection")
        startPing(session)
        return session.send("Hello")
    }

    // 其他方法省略...

    private void startPing(WebSocketSession session) {
        ScheduledExecutorService exec = Executors.newSingleThreadScheduledExecutor()
        exec.scheduleAtFixedRate({
            println("pinging..")
            // 手动订阅并添加背压策略
            session.send("Ping...")
                .onBackpressureBuffer()
                .subscribe(
                    { /* 发送成功 */ },
                    { error -> error.printStackTrace() }
                )
        }, 0, 1, TimeUnit.SECONDS)
    }
}
关键注意点
  1. 必须订阅Publisher:session.send()返回的是冷Observable/Flowable,只有被订阅时才会执行发送逻辑——这也是你之前直接调用session.send()但客户端收不到消息的核心原因。
  2. 处理背压:添加onBackpressureBuffer()或其他背压策略(比如onBackpressureDrop()),避免WebSocket发送队列满时抛出异常。
  3. 线程安全:确保WebSocketSession的操作在合适的线程执行,用subscribeOn(Schedulers.io())可以避免阻塞Netty的事件循环线程。
  4. @OnClose的消息问题:你观察到的@OnClose里send()消息无法送达是正常的,因为此时WebSocket会话已经处于关闭状态,无法再向客户端发送消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 16:39:10