You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多线程任务按累计执行量统一暂停与恢复的实现可行性问询

这个需求完全可以实现,先梳理下你现有代码的问题,再给你正确的实现方式

现有代码的明显问题

  • 变量名混乱:构造函数里this.totalRecords = totalRecords;是笔误,类成员变量应该和传入的totalLogsPrinted对应
  • 暂停逻辑不符合要求:当前代码是每个线程跑完200条才更新计数,而且只暂停当前线程,根本做不到所有线程在累计250条时集体暂停
  • 计数时机错误:你是等整个线程的任务全部完成才去加计数,不是每输出一条就统计,没法精准控制每250条暂停的节点

正确实现思路

要实现每累计250条就让所有线程停10秒,核心解决两个问题:

  1. 精准统计全局完成的任务数,每完成一条就更新计数
  2. 让所有线程在到达暂停节点时互相等待,等全部线程到齐后,统一暂停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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 10:57:52