如何将CompletableFuture<Stream<T>>转为Stream<CompletableFuture<T>>及实现异步任务链
异步任务链展开解决方案
看起来你想实现一个异步任务链:从每个Item生成批量ScrappingResult,再把每个ScrappingResult的视频片段转换成TrimmingResult的异步任务,最终得到所有Trimming任务的Future流。我们来一步步解决你的问题:
核心问题拆解
你面临两个关键问题:
- 如何把
Stream<CompletableFuture<List<ScrappingResult>>>展开为Stream<CompletableFuture<ScrappingResult>> - 如何将每个
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是异步的,无法直接转换,只能选择两种方式:
- 阻塞方式:等待Future完成,取出Stream后包装每个元素为Future:
CompletableFuture<Stream<T>> streamFuture = ...; Stream<CompletableFuture<T>> futureStream = streamFuture.join() .map(CompletableFuture::completedFuture); - 异步方式:用回调将元素的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
相关产品推荐
相关产品推荐

