如何将返回CompletableFuture<Publisher>的Java客户端转为Flux并解决阻塞问题?
问题描述
我正在使用一个返回CompletableFuture<? extends Publisher<?>>的Java客户端,想要将其转换为Flux并实时发射值。由于客户端是无限连接,我怀疑问题出在CompletableFuture永远不会完成?如果是这样,该如何将这个值流转换为Flux并实时发射?
客户端的client.streamQuery()返回CompletableFuture<? extends Publisher<?>>。
下面是我的findAll()方法,它能正常打印"received":
@Override public Flux<Person> findAll() { return Flux.from(client.streamQuery("select * from deal_change emit changes;").join()).flatMap(row -> { System.out.println("received"); try { return Flux.just(mapper.readValue(row.asObject().toJsonString(), Person.class)); } catch (JsonProcessingException e) { return Flux.empty(); } }); }
但调用该方法的代码却没有任何输出,也无法返回结果:
var repo = new PersonRepository(client); var person = repo.findAll().map(p -> { System.out.println("test"); return p; }).blockFirst(); System.out.println(person); System.out.println("end");
解决方案与原因分析
核心问题
当前代码存在两个关键问题:
- 同步阻塞破坏异步链:用
join()同步等待CompletableFuture完成,会阻塞当前线程,打乱Reactor的异步订阅上下文,导致后续流操作无法正确接收数据。 - 错误处理静默吞掉元素:JSON转换失败时返回
Flux.empty(),会直接吞掉当前元素;如果第一个元素就转换失败,blockFirst()会一直等待永远不会出现的元素。
修改后的代码
@Override public Flux<Person> findAll() { // 异步处理CompletableFuture,避免同步阻塞 return Mono.fromFuture(client.streamQuery("select * from deal_change emit changes;")) // 将Publisher转换为Flux .flatMapMany(publisher -> Flux.from(publisher)) // 单个元素转换用map更简洁,替代原flatMap+Flux.just .map(row -> { System.out.println("received"); try { return mapper.readValue(row.asObject().toJsonString(), Person.class); } catch (JsonProcessingException e) { // 抛出异常便于排查,而非静默吞掉元素 throw new RuntimeException("Failed to parse Person from row", e); } }) // 可选:添加错误处理,避免流中断 .onErrorResume(e -> { System.err.println("Processing error: " + e.getMessage()); return Flux.empty(); }); }
关键优化点说明
- 异步处理Future:
Mono.fromFuture会异步监听CompletableFuture的完成事件,完全融入Reactor的异步流模型,不会阻塞线程。 - flatMapMany转换Publisher:当Future完成并返回Publisher后,
flatMapMany将其转换为Flux,实现实时发射元素的需求。 - 错误处理优化:抛出异常可以快速定位转换失败问题,后续的
onErrorResume可选择恢复流或记录日志,避免静默失败。 - map替代flatMap:每个row对应一个Person对象,用
map做同步转换更简洁,无需用Flux.just包装单个元素。
调用端补充建议
如果调用端用blockFirst(),对于无限流来说,它会一直等待第一个元素。可以添加超时时间避免无限阻塞:
var person = repo.findAll() .map(p -> { System.out.println("test"); return p; }) .blockFirst(Duration.ofSeconds(10)); // 设置10秒超时
内容的提问来源于stack exchange,提问作者andresmonc
相关产品推荐
相关产品推荐

