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

多线程读取文件未并行执行问题排查与代码修正请求

问题排查与解决方案

首先咱们直接说核心问题:你只提交了一个Callable任务到线程池,这个任务内部是串行遍历所有的DSection和LView处理文件读取,线程池里的其他线程根本没机会干活,自然是逐个执行,完全没用到多线程并行的优势。

具体问题拆解

  1. 单任务串行执行:整个MultiThreadedFileReadForDashboard是一个单独的Callable,里面嵌套循环处理所有文件,线程池的多线程能力完全没发挥。
  2. 线程池等待方式低效:while (!executor.isTerminated()) {}是忙等,会持续占用CPU资源,应该用awaitTermination实现优雅等待。
  3. 泛型缺失与变量错误:Callable没有指定泛型,导致强制转换不安全;代码里size变量未定义,return dSectionList应该是return dashboardSectionList(构造函数里赋值的是这个变量)。
  4. SFTP连接线程安全问题:ChannelSftp不是线程安全的,多个线程同时使用会出现并发异常,要注意每个任务单独获取连接或者用连接池管理。

修正后的代码实现

1. 拆分独立的文件读取任务

把每个LView的文件读取逻辑拆成单独的Callable,让线程池可以并行处理多个文件:

// 单个文件读取任务,负责处理一个LView
public class FileReadTask implements Callable<LView> {
    private final LView lView;
    private final ChannelSftp sftpChannel; // 注意:若ChannelSftp非线程安全,需给每个任务传新连接
    private final CustomQueryImpl customQuery;

    public FileReadTask(LView lView, ChannelSftp sftpChannel, CustomQueryImpl customQuery) {
        this.lView = lView;
        this.sftpChannel = sftpChannel;
        this.customQuery = customQuery;
    }

    @Override
    public LView call() throws Exception {
        try {
            int userQueryId = Integer.parseInt(lView.getUserQueryId());
            String outputFileName = customQuery.fetchTableInfo(userQueryId);
            if (outputFileName != null && !outputFileName.isBlank()) {
                String data = readFiles(outputFileName);
                lView.setData(data);
            } else {
                lView.setData("No File is present");
            }
        } catch (Exception e) {
            e.printStackTrace();
            lView.setData("Error reading file: " + e.getMessage());
        }
        return lView;
    }

    private String readFiles(String outputFileName) throws Exception {
        StringBuilder inputData = new StringBuilder();
        // 使用try-with-resources自动关闭流,避免资源泄漏
        try (InputStream in = sftpChannel.get(outputFileName);
             BufferedReader br = new BufferedReader(new InputStreamReader(in, "UTF-8"))) {
            String line;
            while ((line = br.readLine()) != null) {
                inputData.append(line).append("\n");
            }
            if (outputFileName.toLowerCase().contains("csv")) {
                JSONArray array = CDL.toJSONArray(inputData.toString());
                return array.toString();
            }
        }
        return "";
    }
}

2. 主线程批量提交任务到线程池

修改Class A的代码,批量提交所有文件读取任务,真正实现并行处理:

class A {
    // 把线程池作为类成员,避免重复创建销毁
    private final ExecutorService executor = getExecuterService();

    private ExecutorService getExecuterService() {
        // IO密集型任务(文件读取)可以适当调大线程数,比如可用核心数*2
        int threadPoolSize = Runtime.getRuntime().availableProcessors() * 2;
        System.out.println("Number of Core Threads: " + threadPoolSize);
        return Executors.newFixedThreadPool(threadPoolSize);
    }

    public List<DSection> processDashboardSections(List<DSection> dashboardSectionList, ChannelSftp sftpChannel, CustomQueryImpl customQuery) throws InterruptedException, ExecutionException {
        List<Future<LView>> futures = new ArrayList<>();

        // 遍历所有节点,提交每个文件读取任务
        for (DSection dSection : dashboardSectionList) {
            List<LView> linkedViewList = new ArrayList<>(dSection.getLinkedViewList());
            for (LView lView : linkedViewList) {
                // 若ChannelSftp非线程安全,这里要为每个任务创建新的连接实例
                futures.add(executor.submit(new FileReadTask(lView, sftpChannel, customQuery)));
            }
        }

        // 等待所有任务完成,处理结果
        for (Future<LView> future : futures) {
            // get()会阻塞直到任务完成,因LView是引用类型,原列表对象已被修改
            future.get();
        }

        // 优雅关闭线程池,设置超时时间避免无限等待
        executor.shutdown();
        if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
            executor.shutdownNow();
        }

        return dashboardSectionList;
    }
}

关键注意事项

  • ChannelSftp线程安全:ChannelSftp实例不能被多线程共享,否则会出现并发异常。建议为每个任务创建新的ChannelSftp连接,或者使用连接池来管理连接。
  • 线程池大小调整:文件读取属于IO密集型任务,线程在等待IO时会处于空闲状态,所以可以把线程池大小设置为可用核心数*2甚至更高,提升资源利用率。
  • 异常处理优化:在任务内部捕获异常并设置错误信息到LView,方便主线程识别失败的任务。
  • 资源管理:用try-with-resources自动关闭文件流,避免资源泄漏问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:45:22