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

如何将CompletableFuture<Stream<T>>转为Stream<CompletableFuture<T>>及实现异步任务链

异步任务链展开解决方案

看起来你想实现一个异步任务链:从每个Item生成批量ScrappingResult,再把每个ScrappingResult的视频片段转换成TrimmingResult的异步任务,最终得到所有Trimming任务的Future流。我们来一步步解决你的问题:

核心问题拆解

你面临两个关键问题:

  1. 如何把Stream<CompletableFuture<List<ScrappingResult>>>展开为Stream<CompletableFuture<ScrappingResult>>
  2. 如何将每个CompletableFuture<ScrappingResult>转换成多个CompletableFuture<TrimmingResult>(对应每个视频片段),并合并成最终的Stream

另外你提到的CompletableFuture<Stream<T>>转Stream<CompletableFuture<T>>,本质是异步结果的同步展开问题——Stream是同步序列,而CompletableFuture是异步的,你无法在Future完成前拿到Stream内的元素,因此需要选择阻塞或异步回调的方式处理。


方案一:异步回调+并发集合(非阻塞,实时处理)

这个方案会在单个Scrapping任务完成后立即启动对应的Trimming任务,不需要等待所有Scrapping任务结束,效率更高:

// 并发队列收集所有Trimming任务的Future
ConcurrentLinkedQueue<CompletableFuture<TrimmingResult>> trimmingFutureQueue = new ConcurrentLinkedQueue<>();

// 第一步:处理每个Item,启动Scrapping异步任务
items.stream()
    .map(item -> CompletableFuture.supplyAsync(() -> new ScrappingTask(item).call()))
    .forEach(listFuture -> {
        // Scrapping任务完成后触发回调
        listFuture.whenCompleteAsync((scrappingResults, error) -> {
            if (error != null) {
                // 处理Scrapping任务异常(比如日志记录)
                error.printStackTrace();
                return;
            }

            // 创建当前Item的输出目录
            Item item = scrappingResults.get(0).getItem();
            Path clipsDir = Paths.get("./" + item.getName() + "/" + item.getTimespan());
            try {
                Files.createDirectories(clipsDir); // 确保目录存在
            } catch (IOException e) {
                e.printStackTrace();
                return;
            }

            AtomicInteger clipIdx = new AtomicInteger();
            // 遍历每个ScrappingResult,处理其视频片段
            scrappingResults.forEach(result -> {
                result.getVideo().getClips().stream()
                    // 为每个片段启动Trimming异步任务
                    .map(clip -> CompletableFuture.supplyAsync(() -> {
                        try {
                            return new TrimmingTask(
                                ffmpegPath,
                                result.getVideo().getVideoUrl(),
                                clip,
                                clipsDir.resolve("clip_" + clipIdx.incrementAndGet() + ".mp3")
                            ).call();
                        } catch (Exception e) {
                            // 处理Trimming任务异常,可返回自定义异常结果或null
                            e.printStackTrace();
                            return null;
                        }
                    }))
                    // 将Trimming任务的Future加入队列
                    .forEach(trimmingFutureQueue::offer);
            });
        });
    });

// 最终得到所有Trimming任务的Future流
Stream<CompletableFuture<TrimmingResult>> trimmingResults = trimmingFutureQueue.stream();

方案二:阻塞式展开(简洁但需等待全部Scrapping任务完成)

如果你能接受等待所有Scrapping任务完成后再启动Trimming任务,代码会更简洁:

// 第一步:收集所有Scrapping任务的Future
List<CompletableFuture<List<ScrappingResult>>> scrappingFutures = items.stream()
    .map(item -> CompletableFuture.supplyAsync(() -> new ScrappingTask(item).call()))
    .collect(Collectors.toList());

// 等待所有Scrapping任务完成(此处会阻塞当前线程)
CompletableFuture.allOf(scrappingFutures.toArray(new CompletableFuture[0])).join();

// 展开所有ScrappingResult,生成单个Result的Future流
Stream<CompletableFuture<ScrappingResult>> scrappingFuturesUnwrapped = scrappingFutures.stream()
    .map(CompletableFuture::join) // 任务已完成,无阻塞
    .flatMap(List::stream)
    .map(CompletableFuture::completedFuture);

// 第二步:生成Trimming任务的Future流
Stream<CompletableFuture<TrimmingResult>> trimmingResults = scrappingFuturesUnwrapped
    .flatMap(resultFuture -> {
        ScrappingResult result = resultFuture.join(); // 任务已完成,无阻塞
        Item item = result.getItem();
        Path clipsDir = Paths.get("./" + item.getName() + "/" + item.getTimespan());
        
        try {
            Files.createDirectories(clipsDir);
        } catch (IOException e) {
            e.printStackTrace();
            return Stream.empty();
        }

        AtomicInteger clipIdx = new AtomicInteger();
        return result.getVideo().getClips().stream()
            .map(clip -> CompletableFuture.supplyAsync(() -> {
                try {
                    return new TrimmingTask(
                        ffmpegPath,
                        result.getVideo().getVideoUrl(),
                        clip,
                        clipsDir.resolve("clip_" + clipIdx.incrementAndGet() + ".mp3")
                    ).call();
                } catch (Exception e) {
                    e.printStackTrace();
                    return null;
                }
            }));
    });

关于CompletableFuture<Stream<T>>转Stream<CompletableFuture<T>>

因为Stream是同步结构,而CompletableFuture是异步的,无法直接转换,只能选择两种方式:

  1. 阻塞方式:等待Future完成,取出Stream后包装每个元素为Future:
    CompletableFuture<Stream<T>> streamFuture = ...;
    Stream<CompletableFuture<T>> futureStream = streamFuture.join()
        .map(CompletableFuture::completedFuture);
    
  2. 异步方式:用回调将元素的Future加入并发集合,再生成Stream:
    ConcurrentLinkedQueue<CompletableFuture<T>> queue = new ConcurrentLinkedQueue<>();
    CompletableFuture<Stream<T>> streamFuture = ...;
    streamFuture.whenCompleteAsync((stream, error) -> {
        if (error != null) {
            error.printStackTrace();
            return;
        }
        stream.map(CompletableFuture::completedFuture)
            .forEach(queue::offer);
    });
    Stream<CompletableFuture<T>> futureStream = queue.stream();
    

注意:异步方式的Stream只有在Future完成后才会有元素,若在Future完成前遍历,可能得到空或部分元素。

内容的提问来源于stack exchange,提问作者Valya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:43:34