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

使用ConnectableFlux实现热流REST接口curl调用返回空数组问题求助

问题根因分析
  • 返回流提前终止:DataService的subscribe方法中,完成对上游热流的订阅后立刻调用了flux.complete(),上游还未发射任何数据,返回给客户端的Flux就已经结束,最终返回空数组。
  • 响应媒体类型配置错误:API接口声明的返回类型为application/json,Spring WebFlux对该媒体类型会收集Flux的所有元素,待流完全结束后才会一次性序列化返回完整JSON数组,无法实现热流逐次推送的效果。
  • 热流管理逻辑缺陷:
    1. 首次创建流时直接调用stream.connect()启动热流,若订阅请求晚于流启动时间到达,会丢失已经发射的历史数据
    2. Flux.push内部仅处理了上游的next事件,未传递完成、异常事件,也未绑定下游取消信号,会存在内存泄漏风险
    3. 内存中存储流的streamsMap没有过期清理逻辑,会随着请求量增长持续占用内存
修复方案

1. 修正订阅逻辑

删除冗余的Flux.push包装和提前调用的complete方法,直接透传上游热流即可:

@Service
class DataService @Autowired constructor(
  private val prv: TestProvider
) {
  fun subscribe(resourceId: String): Flux<QData> {
    return prv.getStream(resourceId)
  }
}

2. 调整返回媒体类型

将接口返回类型改为支持流式推送的格式,可选text/event-stream(SSE标准)或者application/x-ndjson(换行分隔JSON流):

@GetMapping(path = ["/subscription/{resourceId"}], produces = [MediaType.TEXT_EVENT_STREAM_VALUE])
fun subscribe(
  @Parameter(description = "The resource id for which quality data is subscribed for", required = true, example = "example",allowEmptyValue = false)
  @PathVariable("resourceId", required = true) @NotEmpty resourceId: String
): Flux<QData>

3. 优化热流生命周期管理

修改TestProvider的热流创建逻辑,使用replay+autoConnect自动管理流的启动和数据缓存,无需手动维护流存储Map和调用connect:

@Service
class TestProvider {
  fun getStream(resourceId: String): Flux<QData> {
    return Flux.create<QData> { sink ->
      for (i in 1..10) {
        sink.next(QData(LocalDateTime.now().toString(), "next"))
        Thread.sleep(500L)
      }
      sink.complete()
    }
    // 配置缓存最近10条数据,供后加入的订阅者获取历史信息,第一个订阅者到达时自动启动流
    .replay(10)
    .autoConnect()
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 01:48:04