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

