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

使用Files.newDirectoryStream并行递归遍历目录时,如何等待ExecutorService完成所有任务?

使用Files.newDirectoryStream并行递归遍历目录时,如何等待ExecutorService完成所有任务?

问题分析

你遇到的痛点非常典型:因为网络存储延迟高,单线程遍历目录效率太低,想用ExecutorService并行处理,但常规的shutdown()+awaitTermination()方案行不通——毕竟遍历过程中会动态提交子目录的遍历任务;CountDownLatch也没法用,因为我们根本没法提前知道总共有多少任务要执行。

我们测试有效的解决方案

核心思路是用原子计数器跟踪活跃的目录遍历任务数,配合锁机制等待计数器归0,以此判断所有遍历任务完成。这个方案既能支持动态提交新任务,又能准确等待所有任务结束,同时还做了线程池的优化避免内存溢出。

完整代码实现

public class FolderWalker {

    private final ExecutorService executorService;
    private final PathProcessor processor;
    private final PathFilter fileFilter;
    private final PathFilter directoryFilter;
    private final ExceptionHandler exceptionHandler;
    private final AtomicInteger activeThreads = new AtomicInteger(0);
    private final Object lock = new Object();
    
    private boolean keepRunning = true;

    public static interface PathProcessor{
        public void process(Path p, DosFileAttributes attribs) throws Exception;
    }
    
    public static interface ExceptionHandler{
        public boolean handle(Path p, Exception e); // 返回true继续遍历
    }
    
    public static interface PathFilter{
        public boolean accept(Path p, DosFileAttributes attribs);
        
        public static PathFilter ACCEPT_ALL = (p, attribs) -> true;
    }

    public FolderWalker(int numThreads, PathFilter directoryFilter, PathFilter fileFilter, PathProcessor processor, ExceptionHandler exceptionHandler) {
        this.executorService = createExecutorService(numThreads);
        this.directoryFilter = directoryFilter;
        this.fileFilter = fileFilter;
        this.processor = processor;
        this.exceptionHandler = exceptionHandler;
    }
    
    private static ExecutorService createExecutorService(int numThreads) {
        
        // 当所有线程都在忙时,让提交任务的线程自己执行任务,避免任务队列无限膨胀导致内存溢出
        return new ThreadPoolExecutor(numThreads, numThreads, 0, TimeUnit.MILLISECONDS, 
                  new SynchronousQueue<>(), 
                  new ThreadPoolExecutor.CallerRunsPolicy());
    }

    public void walkDirectory(Path directory) {
        activeThreads.incrementAndGet();
        executorService.submit(() -> walkDirectoryInternal(directory));
    }
    
    private void walkDirectoryInternal(Path directory) {
        int activeThreadCount = activeThreads.get();
        System.out.println(Thread.currentThread() + " - Before: Looking in " + directory + " - Active folders: " + activeThreadCount);
        
        try {
            try (DirectoryStream<Path> stream = Files.newDirectoryStream(directory)) {
                for (Path entry : stream) {
                    if (Thread.interrupted()) return;
                    if (!keepRunning) return;
                    
                    try {
                        DosFileAttributes attribs =  getBasicFileAttributes(entry);
                        if (attribs.isDirectory() && directoryFilter.accept(entry, attribs)) {
                            activeThreads.incrementAndGet();
                            executorService.submit(() -> walkDirectoryInternal(entry));
                        } else if (fileFilter.accept(entry, attribs)) {
                            processor.process(entry, attribs);
                        }
                    } catch (Exception e) {
                        if (!exceptionHandler.handle(entry, e))
                            keepRunning = false;
                    }
                }
            } catch (Exception e) {
                if (!exceptionHandler.handle(directory, e))
                    keepRunning = false;
            }
        } finally {
            activeThreadCount = activeThreads.decrementAndGet();
            System.out.println(Thread.currentThread() + " - After: Looking in " + directory + " - Active folders: " + activeThreadCount);
            if (activeThreadCount == 0) {
                synchronized(lock) {
                    lock.notifyAll();
                }
            }
        }           
            
    }

    private DosFileAttributes getBasicFileAttributes(Path p) throws IOException {
        return Files.getFileAttributeView(p, DosFileAttributeView.class).readAttributes();
    }
    
    public void await() throws Exception {
        synchronized(lock) {
            while (activeThreads.get() != 0)
                lock.wait();
        }
        
        executorService.shutdown();
        System.out.println("线程池已关闭");
        try {
            if (!executorService.awaitTermination(60, TimeUnit.SECONDS)) {
                executorService.shutdownNow();
            }
        } catch (InterruptedException e) {
            executorService.shutdownNow();
        }
        System.out.println("线程池已终止");
        
    }

}

关键逻辑说明

  • 原子计数器activeThreads:每次提交目录遍历任务前,先调用incrementAndGet()递增计数;任务执行完成后(finally块中)调用decrementAndGet()递减计数,精准跟踪当前正在执行的目录遍历任务数。
  • 锁等待机制:当计数器降到0时,说明所有目录遍历任务都已完成,此时通过lock.notifyAll()唤醒等待的主线程;await()方法则通过lock.wait()阻塞,直到计数器归0。
  • 线程池优化:采用SynchronousQueue+CallerRunsPolicy的组合,避免任务队列无限膨胀导致内存溢出——当所有线程都在忙时,新任务会由提交任务的线程自行处理,这在遍历海量目录时尤为重要。
  • 异常与中断控制:通过keepRunning标志控制遍历是否继续,同时处理了线程中断的情况,保证遍历过程能优雅停止。

补充建议

  • 如果想替换锁等待的逻辑,也可以考虑用Phaser(Java 7+),不过原子计数器+锁的方式更直观易懂。
  • 异常处理逻辑可以根据业务需求调整,比如遇到致命错误时立即终止所有遍历任务。
  • 线程池大小建议根据网络带宽和存储性能调整,不要设置过大,避免过多线程竞争导致性能下降。

备注:内容来源于stack exchange,提问作者Kevin Day

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 19:12:58