多线程任务按累计执行量统一暂停与恢复的实现可行性问询
这个需求完全可以实现,先梳理下你现有代码的问题,再给你正确的实现方式
现有代码的明显问题
- 变量名混乱:构造函数里
this.totalRecords = totalRecords;是笔误,类成员变量应该和传入的totalLogsPrinted对应 - 暂停逻辑不符合要求:当前代码是每个线程跑完200条才更新计数,而且只暂停当前线程,根本做不到所有线程在累计250条时集体暂停
- 计数时机错误:你是等整个线程的任务全部完成才去加计数,不是每输出一条就统计,没法精准控制每250条暂停的节点
正确实现思路
要实现每累计250条就让所有线程停10秒,核心解决两个问题:
- 精准统计全局完成的任务数,每完成一条就更新计数
- 让所有线程在到达暂停节点时互相等待,等全部线程到齐后,统一暂停10秒再继续
这里用CyclicBarrier最合适,它能让一组线程互相等待,直到所有线程都到达同步点,再一起执行后续逻辑。
修正后的完整代码
import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class LogProducer { private static final Logger log = LoggerFactory.getLogger(LogProducer.class); public void produceLog() { int threadCount = 5; int totalTasks = 1000; int pauseThreshold = 250; AtomicInteger totalCompletedTasks = new AtomicInteger(0); // 初始化CyclicBarrier:当threadCount个线程都到达屏障时,执行暂停逻辑 CyclicBarrier pauseBarrier = new CyclicBarrier(threadCount, () -> { try { log.info("累计完成{}条任务,所有线程暂停10秒", totalCompletedTasks.get()); TimeUnit.SECONDS.sleep(10); log.info("暂停结束,所有线程恢复执行"); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } }); ExecutorService executor = Executors.newFixedThreadPool(threadCount); int tasksPerThread = totalTasks / threadCount; for (int i = 0; i < threadCount; i++) { executor.execute(new TaskWorker(tasksPerThread, pauseThreshold, totalCompletedTasks, pauseBarrier)); } executor.shutdown(); try { // 等待所有任务完成,超时时间设为1小时足够覆盖任务+暂停时间 if (!executor.awaitTermination(1, TimeUnit.HOURS)) { executor.shutdownNow(); } } catch (InterruptedException e) { executor.shutdownNow(); Thread.currentThread().interrupt(); } } static class TaskWorker implements Runnable { private final int tasksToRun; private final int pauseThreshold; private final AtomicInteger totalCompleted; private final CyclicBarrier barrier; public TaskWorker(int tasksToRun, int pauseThreshold, AtomicInteger totalCompleted, CyclicBarrier barrier) { this.tasksToRun = tasksToRun; this.pauseThreshold = pauseThreshold; this.totalCompleted = totalCompleted; this.barrier = barrier; } @Override public void run() { try { for (int i = 0; i < tasksToRun; i++) { // 这里替换成你的耗时业务任务 log.info("Hello World | 线程ID: {}", Thread.currentThread().getId()); // 每完成一条就更新全局计数 int completed = totalCompleted.incrementAndGet(); // 判断是否到达暂停节点(排除最后1000条的情况,避免不必要的暂停) if (completed % pauseThreshold == 0 && completed != 1000) { // 等待所有线程都到达这个暂停点 barrier.await(); } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.error("线程被中断", e); } catch (BrokenBarrierException e) { log.error("同步屏障异常", e); } } } public static void main(String[] args) { new LogProducer().produceLog(); } }
代码关键点说明
- 原子计数:用
AtomicInteger保证多线程下计数的准确性,每完成一条任务立即更新 - CyclicBarrier的作用:当所有线程都完成当前阶段任务(累计到250/500/750条),触发屏障回调函数,让所有线程一起暂停10秒,之后继续执行剩余任务
- 线程池管理:用
shutdown()和awaitTermination()确保所有任务完成后正确关闭线程池,避免资源泄漏
内容的提问来源于stack exchange,提问作者user1609
相关产品推荐
相关产品推荐

