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

基于Reactor/WebFlux实现无资源泄漏的响应式HTTP新闻流推送

嘿,我完全懂你的需求——想用流式JSON(application/stream+json)推送新闻事件,不想碰WebSocket,还得解决烦人的资源泄漏问题对吧?刚好在Spring WebFlux里做这个有标准的、安全的做法,我来给你一步步讲清楚。

核心思路:用Sinks实现安全的事件推送

Spring WebFlux推荐用Sinks来创建可推送的Flux,它比老的EmitterProcessor更安全,内置了背压处理和订阅者生命周期管理,能有效避免资源泄漏。

完整代码实现

首先定义你的News实体(用record更简洁):

import java.time.Instant;

public record News(String id, String content, Instant timestamp) {}

然后是RestController的实现,包含流式接口和新闻发布入口:

import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;

@RestController
public class NewsStreamController {

    // 多播场景:多个客户端共享同一份新闻流(适合广播)
    // 如果是单播(每个客户端独立流),改用Sinks.many().unicast().onBackpressureBuffer()
    private final Sinks.Many<News> newsSink = Sinks.many().multicast().onBackpressureBuffer();

    // 对外暴露的流式接口,返回application/stream+json
    @GetMapping(value = "/news/stream", produces = MediaType.APPLICATION_STREAM_JSON_VALUE)
    public Flux<News> streamNews() {
        return newsSink.asFlux()
                // 客户端断开连接时触发,做清理工作,避免资源泄漏
                .doOnCancel(() -> {
                    System.out.println("客户端断开,已取消新闻流订阅");
                    // 这里可以加自定义清理逻辑,比如移除订阅者相关资源
                })
                // 处理流中的错误,避免错误扩散导致资源泄漏
                .doOnError(error -> {
                    newsSink.tryEmitError(error);
                    System.err.println("新闻流发生错误: " + error.getMessage());
                });
    }

    // 模拟外部发布新闻的方法(比如从其他服务、MQ、用户接口调用)
    public void publishNews(News news) {
        // 安全发送事件,自动处理背压和订阅者状态
        Sinks.EmitResult result = newsSink.tryEmitNext(news);
        if (result.isFailure()) {
            // 处理发送失败的情况(比如无订阅者、缓冲区满)
            System.err.println("新闻发布失败: " + result.getReason());
        }
    }
}

为什么这样不会泄漏资源?

  • 订阅者生命周期管理:Sinks会自动跟踪订阅者,当客户端断开连接(触发cancel),doOnCancel会执行清理,避免持有无用的订阅者引用。
  • 背压处理:onBackpressureBuffer会合理管理消息缓冲区,不会无限制占用内存;你也可以根据需求调整策略(比如丢弃旧消息)。
  • 线程安全:Sinks本身是线程安全的,publishNews可以在任何线程调用,不用担心并发问题。

常见的资源泄漏坑点(你可能踩过的)

  • 手动创建Flux时没处理cancel信号:比如用Flux.create但没在Sink的onCancel回调里清理资源。
  • 使用了未正确管理的Processor:老的EmitterProcessor如果没处理订阅取消,容易残留订阅者引用。
  • 忽略背压:导致消息缓冲区无限增长,占用内存直至OOM。

测试一下

启动应用后,用curl调用接口就能看到流式输出:

curl http://localhost:8080/news/stream

然后调用publishNews方法发布新闻,curl里会实时收到JSON消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:11:15