Spring Boot WebFlux:如何持续查询数据库最新记录并输出Flux
解决方案
你的代码核心问题在于使用了map操作符处理返回Mono的数据库查询方法——map只能同步转换元素,无法处理异步的Mono结果,最终会导致Flux发布的是Mono对象而非实际的字段值。正确的做法是使用flatMap操作符,它可以订阅异步的Mono并将其结果整合到Flux流中。
正确实现代码
@RestController public class WebFluxController { private final ReactiveDocumentRepository reactiveDocumentRepository; // 构造注入替代@Autowired(更符合Spring最佳实践) public WebFluxController(ReactiveDocumentRepository reactiveDocumentRepository) { this.reactiveDocumentRepository = reactiveDocumentRepository; } @CrossOrigin @GetMapping(value = "/documents") public Flux<String> getDocuments() { return Flux.interval(Duration.ofSeconds(1)) // 用flatMap处理异步Mono .flatMap(tick -> reactiveDocumentRepository.findLastDocument() // 提取description字段 .map(Document::getDescription) // 处理无最新记录的情况 .defaultIfEmpty("暂无最新文档") // 捕获数据库查询异常,返回友好提示 .onErrorResume(e -> Mono.just("查询最新文档失败:" + e.getMessage()))); } }
关键细节说明
- 使用
flatMap替代map:Flux.interval每秒发出一个时间戳,flatMap会为每个时间戳触发一次数据库查询,订阅findLastDocument()返回的Mono,并将其结果(文档的description字段)作为Flux的元素发布。 - 处理空结果:用
defaultIfEmpty在数据库无最新记录时返回默认值,避免流中断。 - 异常处理:
onErrorResume捕获数据库查询过程中的异常,返回错误提示,保证Flux流不会因异常终止。 - 构造注入:替代
@Autowired,提升代码的可测试性和依赖透明度,符合Spring的最佳实践。
内容的提问来源于stack exchange,提问作者DariusNica
相关产品推荐
相关产品推荐

