多线程任务协调与线程池管理实践
在处理大规模数据导出或并行查询等场景中,合理使用线程池和同步工具(如 CountDownLatch 和 CompletableFuture)可显著提升系统吞吐量与响应速度。
1. 使用 CountDownLatch 协调多线程任务完成
以下示例展示了如何通过 CountDownLatch 等待多个异步任务全部执行完毕后再进行后续操作(如生成 Excel 文件):
ExecutorService executor = new ThreadPoolExecutor(
10, 10, 0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>()
);
List<Map<String, Object>> rawData = reader.readAll();
CountDownLatch latch = new CountDownLatch(rawData.size());
rawData.forEach(record -> {
executor.execute(() -> {
try {
// 执行业务逻辑:查询视频、转码状态等
processRecord(record);
} finally {
latch.countDown(); // 确保计数器递减
}
});
});
try {
latch.await(); // 阻塞直到所有任务完成
executor.shutdown();
// 导出结果到 Excel
ExcelWriter writer = ExcelUtil.getWriter(new File("output.xlsx"));
writer.write(rawData);
writer.close();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("任务等待被中断", e);
}
2. 获取异步任务返回值:submit + Future
当需要获取子线程的计算结果时,应使用 submit(Callable) 而非 execute(Runnable):
List<Future<Integer>> futures = new ArrayList<>();
CountDownLatch doneSignal = new CountDownLatch(dataList.size());
for (Integer item : dataList) {
Future<Integer> future = executor.submit(() -> {
try {
return processItem(item); // 返回计算结果
} finally {
doneSignal.countDown();
}
});
futures.add(future);
}
doneSignal.await();
executor.shutdown();
// 收集结果
List<Integer> results = new ArrayList<>();
for (Future<Integer> f : futures) {
results.add(f.get()); // 可能抛出 ExecutionException
}
3. 使用 CompletableFuture 实现更灵活的并行组合
对于多个独立查询任务,CompletableFuture 提供了更简洁的异步编排能力:
CompletableFuture<List<AppGlobalSearch>> appFuture =
CompletableFuture.supplyAsync(() -> globalSearchApp(qo), searchExecutor);
CompletableFuture<List<IssueDto>> issueFuture =
CompletableFuture.supplyAsync(() -> searchIssues(...), searchExecutor);
// 等待所有任务完成(带超时)
CompletableFuture.allOf(appFuture, issueFuture)
.get(3, TimeUnit.SECONDS);
// 安全获取结果(即使部分失败)
List<AppGlobalSearch> apps = safeGet(appFuture);
List<IssueDto> issues = safeGet(issueFuture);
private <T> List<T> safeGet(CompletableFuture<List<T>> future) {
try {
return future.getNow(Collections.emptyList());
} catch (Exception e) {
return Collections.emptyList();
}
}
4. 线程池类型与关闭策略
- newFixedThreadPool(n):固定大小线程池,适用于负载较重的服务器。
- newCachedThreadPool():弹性线程池,适合执行大量短期异步任务。
- shutdown():平滑关闭,已提交任务继续执行,但不再接受新任务。
- shutdownNow():尝试立即停止所有正在执行的任务,返回未执行的任务列表。
建议在线程池使用完毕后调用 shutdown(),并在必要时配合 awaitTermination() 等待其真正终止。
5. CountDownLatch 核心机制
CountDownLatch(int count):初始化计数器。countDown():每次调用使计数减一。await():阻塞当前线程,直到计数归零。await(timeout, unit):支持超时的等待,避免永久阻塞。
典型用途:主线程等待 N 个子任务全部完成后才继续执行后续逻辑。
6. 线程池监控指标
可通过 ThreadPoolExecutor 的方法获取运行时状态:
int active = ((ThreadPoolExecutor) pool).getActiveCount(); // 活跃线程数
long totalTasks = ((ThreadPoolExecutor) pool).getTaskCount(); // 总提交任务数
long completed = ((ThreadPoolExecutor) pool).getCompletedTaskCount(); // 已完成任务数
int poolSize = ((ThreadPoolExecutor) pool).getPoolSize(); // 当前线程池大小
int maxEver = ((ThreadPoolExecutor) pool).getLargestPoolSize(); // 历史最大线程数
