当前位置:首页 > 随笔 > 正文内容

多线程任务协调与线程池管理实践

访客 随笔 2026年8月7日 1

在处理大规模数据导出或并行查询等场景中,合理使用线程池和同步工具(如 CountDownLatchCompletableFuture)可显著提升系统吞吐量与响应速度。

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(); // 历史最大线程数

相关文章

可以按小时收费的VPS

很多 VPS 提供商都支持 按小时计费(hourly billing),想短期试用 / 临时搭建节点、测试网络、短期项目等场景非常合适。下面是当前最主流且靠谱的按小时 VPS 选项,分别按不同需求场景整理: 1. Vultr(全球节点,包括日本) 按小时计费 可选机房:东京 / 大阪 / 洛杉矶 / 法兰克福 / 伦敦 … 支持 PayPal(部分情况),但更常用信用卡/PayPal+卡价格参考$...

在 iPhone 上下载国外App

地区/国家限制App Store 会根据 Apple ID 的国家或地区限制应用下载。如果你的 Apple ID 绑定的是中国大陆,就可能无法下载 OpenAI 官方的 ChatGPT 应用,因为它在大陆 App Store 不上架。解决办法:换成美国、加拿大、香港等地区的 Apple ID。或者在现有 Apple ID 上更改地区。注册一个国外 Apple ID(推荐)比如注册 美国区 Appl...

Node.js 中的异步编程:回调与 Promise

Node.js 是一个基于 JavaScript 构建的单线程、非阻塞运行环境,它通过异步编程机制来高效处理多个操作。在执行如文件读取、API 请求或数据库查询等任务时,Node.js 不会等待这些操作完成,而是使用回调函数和 Promise 来避免阻塞主线程。 回调方式实现异步 那么当异步操作完成后,Node.js 如何知道接下来要做什么呢?这就要用到 回调函数(callback)。 回调本质上...

Selenium自动化测试入门指南

Selenium自动化测试入门指南

什么是自动化测试? 自动化测试是指利用软件工具自动执行测试用例,模拟用户操作,如打开网页、点击链接、输入文本等,并验证结果是否符合预期。 其主要优点包括: 大幅减少人工成本 测试速度快 可以在非工作时间运行 支持持续集成和交付 然而,它也存在一些局限性,例如开发成本较高、不适合快速变化的项目、依赖稳定的UI界面等。 自动化测试的应用条件 适合引入自动化测试的情况包括: 手动测试耗时且需要大量...

MariaDB Galera集群故障快速恢复指南

OpenStack控制节点采用三节点MariaDB Galera集群架构。当数据库集群因故障重启时,有时会出现Galera集群无法正常启动的问题。虽然有多种方法可以恢复数据库服务,但如何实现快速启动同时确保数据完整性呢? 通过分析日志发现,MariaDB Galera集群节点宕机时会在日志中输出以下信息: [Note] WSREP: 新集群视图:全局状态: 874d8e7e-5980-11e8-8...

Android 中 EventBus 的通信机制与实现原理深度解析

EventBus 核心设计思想 EventBus 是一个基于观察者模式的事件总线框架,广泛应用于 Android 平台以实现组件解耦。它通过中心化的消息分发机制,使不同层级、不同线程的对象能够以"发布-订阅"方式通信,避免了传统接口回调或广播带来的强依赖问题。 核心角色说明 事件(Event):任意 Java 对象,作为数据载体,如网络状态变更通知、用户登录信息等。 发布者(Publi...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。