多线程读取文件未并行执行问题排查与代码修正请求
问题排查与解决方案
首先咱们直接说核心问题:你只提交了一个Callable任务到线程池,这个任务内部是串行遍历所有的DSection和LView处理文件读取,线程池里的其他线程根本没机会干活,自然是逐个执行,完全没用到多线程并行的优势。
具体问题拆解
- 单任务串行执行:整个
MultiThreadedFileReadForDashboard是一个单独的Callable,里面嵌套循环处理所有文件,线程池的多线程能力完全没发挥。 - 线程池等待方式低效:
while (!executor.isTerminated()) {}是忙等,会持续占用CPU资源,应该用awaitTermination实现优雅等待。 - 泛型缺失与变量错误:Callable没有指定泛型,导致强制转换不安全;代码里
size变量未定义,return dSectionList应该是return dashboardSectionList(构造函数里赋值的是这个变量)。 - 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
相关产品推荐
相关产品推荐

