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

如何将返回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();
            });
}

关键优化点说明

  1. 异步处理Future:Mono.fromFuture会异步监听CompletableFuture的完成事件,完全融入Reactor的异步流模型,不会阻塞线程。
  2. flatMapMany转换Publisher:当Future完成并返回Publisher后,flatMapMany将其转换为Flux,实现实时发射元素的需求。
  3. 错误处理优化:抛出异常可以快速定位转换失败问题,后续的onErrorResume可选择恢复流或记录日志,避免静默失败。
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 15:48:31