如何为大文本文件解析任务应用ExecutorService多线程处理
基于ExecutorService实现多线程处理大文件的改造方案
核心思路
保留现有FileParser逐行解析的逻辑,将单线程同步校验API的环节拆分为并行任务,通过线程池执行,最后统一汇总结果生成报告。重点解决线程安全、任务调度、资源释放三个核心问题。
步骤1:配置线程池
根据CPU核心数和API接口的并发承载能力,自定义线程池参数(比Executors默认实现更可控):
// 可根据实际场景调整参数 private static final ExecutorService executor = new ThreadPoolExecutor( 4, // 核心线程数:保持存活的最小线程数 8, // 最大线程数:峰值并发数 60L, TimeUnit.SECONDS, // 空闲线程存活时间 new LinkedBlockingQueue<>(1000), // 任务队列:缓冲待执行的校验任务 Executors.defaultThreadFactory(), new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时,由提交任务的线程直接执行,避免丢任务 );
步骤2:改造解析-校验流程
将原来的单线程同步校验,改为解析后提交任务到线程池,等待所有任务完成后生成报告。注意使用线程安全容器存储结果:
// 线程安全的结果容器 private final ConcurrentMap<FileInfo, Boolean> checkResults = new ConcurrentHashMap<>(); public void processFile(InputStream inputStream, File outputFile) { try (BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream))) { String line; List<Future<?>> futures = new ArrayList<>(); // 逐行解析,提交异步校验任务 while ((line = reader.readLine()) != null) { FileInfo fileInfo = new FileParser().parseLine(line); Future<?> future = executor.submit(() -> { boolean exists = false; try { exists = apiService.checkExists(fileInfo); } catch (Exception e) { // 单独捕获任务异常,避免线程池线程因未捕获异常终止 System.err.printf("校验FileInfo失败: %s, 原因: %s%n", fileInfo, e.getMessage()); } checkResults.put(fileInfo, exists); }); futures.add(future); } // 等待所有校验任务执行完毕 for (Future<?> future : futures) { try { future.get(); // 阻塞等待任务完成 } catch (InterruptedException | ExecutionException e) { System.err.printf("任务执行异常: %s%n", e.getMessage()); } } } catch (IOException e) { System.err.printf("读取输入文件失败: %s%n", e.getMessage()); } finally { // 关闭线程池:不再接受新任务,等待已有任务执行完成 executor.shutdown(); try { if (!executor.awaitTermination(1, TimeUnit.HOURS)) { executor.shutdownNow(); // 超时未完成则强制关闭 } } catch (InterruptedException e) { executor.shutdownNow(); } } // 所有任务完成后生成报告 new ReportGenerator().generate(checkResults, outputFile); }
步骤3:关键细节优化
ApiService线程安全校验
如果ApiService是无状态的(仅调用HTTP接口、无共享变量),无需额外处理;若存在共享状态,需添加synchronized锁或使用线程安全组件。超大型文件内存优化
若文件过大,ConcurrentMap可能占用过多内存,可改为边校验边写报告:- 用
BlockingQueue传递校验结果 - 单独启动一个线程负责写入报告,避免阻塞校验线程
// 结果传递队列 private final BlockingQueue<Result> resultQueue = new LinkedBlockingQueue<>(); public void processFile(InputStream inputStream, File outputFile) { // 启动独立的写报告线程 Thread reportWriter = new Thread(() -> { try (BufferedWriter writer = new BufferedWriter(new FileWriter(outputFile))) { Result result; while ((result = resultQueue.take()) != null) { if (result.fileInfo == null) break; // 任务结束信号 // 按业务格式写入报告 writer.write(String.format("%s,%b", result.fileInfo.getId(), result.exists)); writer.newLine(); } } catch (IOException | InterruptedException e) { System.err.printf("写入报告失败: %s%n", e.getMessage()); } }); reportWriter.start(); // 解析+提交校验任务 try (BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream))) { String line; List<Future<?>> futures = new ArrayList<>(); while ((line = reader.readLine()) != null) { FileInfo fileInfo = new FileParser().parseLine(line); futures.add(executor.submit(() -> { boolean exists = false; try { exists = apiService.checkExists(fileInfo); } catch (Exception e) { System.err.printf("校验失败: %s%n", fileInfo); } resultQueue.put(new Result(fileInfo, exists)); })); } // 等待所有校验任务完成 for (Future<?> future : futures) future.get(); resultQueue.put(new Result(null, false)); // 发送结束信号 } catch (Exception e) { System.err.printf("处理文件失败: %s%n", e.getMessage()); } finally { executor.shutdown(); try { reportWriter.join(); // 等待写报告线程结束 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } // 结果封装类 private static class Result { FileInfo fileInfo; boolean exists; Result(FileInfo fileInfo, boolean exists) { this.fileInfo = fileInfo; this.exists = exists; } }- 用
API超时控制
给ApiService的校验接口添加超时,避免单个任务长期阻塞线程池:public boolean checkExists(FileInfo fileInfo) throws IOException { HttpRequest request = HttpRequest.newBuilder() .uri(URI.create("你的校验API地址")) .timeout(Duration.ofSeconds(10)) // 设置10秒超时 .POST(HttpRequest.BodyPublishers.ofString(fileInfo.toJson())) .build(); HttpResponse<String> response = HttpClient.newHttpClient().send(request, HttpResponse.BodyHandlers.ofString()); return response.statusCode() == 200 && Boolean.parseBoolean(response.body()); }
内容的提问来源于stack exchange,提问作者user19772950
相关产品推荐
相关产品推荐

