为什么我的Spring WebFlux控制器仅在第一次请求时返回数据?
问题根因
你当前的方案存在两个核心问题:
- 你使用的
Sinks.many().multicast()是多播模式,默认不会缓存历史发送的元素,订阅者仅能收到订阅完成后推送的新元素。你的接口逻辑是先把请求推送到请求Sink,再订阅响应Sink,当externalService处理速度较快时,响应结果会在你订阅响应Sink之前就被推送出去,后续订阅自然拿不到任何数据,就会返回空响应。 - 全局共用同一套请求/响应Sink,没有做请求和响应的绑定逻辑,并发请求场景下会出现A请求的结果被B请求拿到的串数据问题,完全不符合业务要求。
优化方案
你的需求完全不需要引入全局Sink来做异步解耦,直接利用WebFlux自带的异步调度和超时能力就能实现,简化后的实现代码如下:
@RequestMapping("/reactive") @RestController class SampleController @Autowired constructor(private val externalService: ExternalService) { @GetMapping("/getUser/{phoneNumber}") fun getUser(@PathVariable phoneNumber: String): Mono<String> { // 把阻塞的外部查询包成Mono,指定弹性线程池运行,不占用WebFlux核心IO线程 val queryTask = Mono.fromCallable { externalService.findByIdOrNull(phoneNumber) } .subscribeOn(Schedulers.boundedElastic()) // 处理完成后无论是否触发超时,都发送邮件通知用户,这里替换成你自己的邮件发送逻辑 .doOnNext { user -> if (user != null) { sendResultEmail(phoneNumber, user) } else { sendErrorEmail(phoneNumber, "未查询到对应用户信息") } } .doOnError { e -> sendErrorEmail(phoneNumber, e.message ?: "查询异常") } // 缓存结果避免重复执行查询 .cache() return queryTask .map { it.toString() } // 20秒内返回结果则直接返回,超时则返回提示文案,后台查询任务会继续运行直到完成后发邮件 .timeout(Duration.ofSeconds(20), Mono.just("your request is under process")) } // 邮件发送方法示例,自行实现 private fun sendResultEmail(phoneNumber: String, user: AppUser) { // 实现发送结果邮件逻辑 } private fun sendErrorEmail(phoneNumber: String, errMsg: String) { // 实现发送异常通知邮件逻辑 } }
如果你确实有使用Sink的场景需求,也应该做到每个请求对应独立的Sinks.one()实例,或者给请求增加唯一标识做响应匹配,不要全局共用多播Sink。
内容的提问来源于stack exchange,提问作者hamid
相关产品推荐
相关产品推荐

