如何重写主线程-工作线程的同步逻辑?
嘿,我来帮你梳理下怎么重写这个主线程与工作线程的同步逻辑~先提个小细节,你代码里的内层循环for(int j=0 ; j = 20 ; j++)是个bug——这里用了赋值运算符=而不是比较运算符,应该改成j < 20,不然这个循环会一直跑下去根本停不住😂。
接下来咱们聊聊几种靠谱的同步方案,适配你代码里“多轮次、主线程+工作线程双向同步”的需求:
方案1:优化现有CyclicBarrier逻辑
你原本的思路是用CyclicBarrier让主线程和所有工作线程在每轮计算的“准备阶段”和“完成阶段”同步,这个方向是对的,但需要把主线程的同步逻辑补全,和工作线程的节奏对齐。
修正后的完整代码:
import java.util.concurrent.CyclicBarrier; public class Test implements Runnable { public int local_counter; public static int global_counter; // 定义工作线程数量,这里示例设为4 public static final int n_threads = 4; // Barrier等待n_threads个工作线程 + 1个主线程 public static CyclicBarrier thread_barrier = new CyclicBarrier(n_threads + 1); @Override public void run() { try { for (int i = 0; i < 100; i++) { // 第一阶段同步:等待主线程发号施令,所有线程一起开始本轮计算 thread_barrier.await(); local_counter = 0; // 修正内层循环条件 for (int j = 0; j < 20; j++) { local_counter++; } // 第二阶段同步:等待所有线程完成本轮计算,主线程可以执行后续操作 thread_barrier.await(); } } catch (Exception e) { e.printStackTrace(); } } public static void main(String[] args) { Thread[] threads = new Thread[n_threads]; // 启动所有工作线程 for (int i = 0; i < n_threads; i++) { threads[i] = new Thread(new Test()); threads[i].start(); } try { for (int i = 0; i < 100; i++) { // 和工作线程同步:主线程就绪,触发本轮计算开始 thread_barrier.await(); // 这里可以添加主线程在每轮的操作,比如统计全局计数 global_counter += 1; // 和工作线程同步:等待所有工作线程完成计算,进入下一轮 thread_barrier.await(); } // 等待所有工作线程执行完毕 for (Thread t : threads) { t.join(); } } catch (Exception e) { e.printStackTrace(); } } }
核心逻辑:每一轮循环里,主线程和所有工作线程在两个节点严格同步——先一起准备好开始计算,再一起完成计算后进入下一轮,完美对齐节奏。
方案2:CountDownLatch + CyclicBarrier组合(更严谨的启动同步)
如果担心工作线程启动时间不一致,导致第一轮同步出现偏差,可以先用CountDownLatch确保所有工作线程都启动完成后,主线程再进入循环:
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; public class Test implements Runnable { public int local_counter; public static int global_counter; public static final int n_threads = 4; // 用来等待所有工作线程启动完成 public static CountDownLatch startLatch = new CountDownLatch(n_threads); // 轮次同步的Barrier public static CyclicBarrier phaseBarrier = new CyclicBarrier(n_threads + 1); @Override public void run() { try { // 通知主线程:当前工作线程已启动 startLatch.countDown(); for (int i = 0; i < 100; i++) { phaseBarrier.await(); local_counter = 0; for (int j = 0; j < 20; j++) { local_counter++; } phaseBarrier.await(); } } catch (Exception e) { e.printStackTrace(); } } public static void main(String[] args) { Thread[] threads = new Thread[n_threads]; for (int i = 0; i < n_threads; i++) { threads[i] = new Thread(new Test()); threads[i].start(); } try { // 先等所有工作线程都启动完毕 startLatch.await(); for (int i = 0; i < 100; i++) { phaseBarrier.await(); global_counter += 1; phaseBarrier.await(); } for (Thread t : threads) { t.join(); } } catch (Exception e) { e.printStackTrace(); } } }
方案3:用Phaser实现更灵活的同步
如果后续可能有动态增减线程的需求,Phaser比CyclicBarrier更灵活,它支持动态注册/注销参与者:
import java.util.concurrent.Phaser; public class Test implements Runnable { public int local_counter; public static int global_counter; public static final int n_threads = 4; // 初始注册n_threads+1个参与者(工作线程+主线程) public static Phaser phaser = new Phaser(n_threads + 1); @Override public void run() { try { for (int i = 0; i < 100; i++) { // 等待所有参与者到达当前阶段 phaser.arriveAndAwaitAdvance(); local_counter = 0; for (int j = 0; j < 20; j++) { local_counter++; } // 再次等待所有参与者完成当前阶段 phaser.arriveAndAwaitAdvance(); } // 工作线程结束,注销自己 phaser.arriveAndDeregister(); } catch (Exception e) { e.printStackTrace(); } } public static void main(String[] args) { Thread[] threads = new Thread[n_threads]; for (int i = 0; i < n_threads; i++) { threads[i] = new Thread(new Test()); threads[i].start(); } try { for (int i = 0; i < 100; i++) { phaser.arriveAndAwaitAdvance(); global_counter += 1; phaser.arriveAndAwaitAdvance(); } // 主线程结束,注销自己 phaser.arriveAndDeregister(); for (Thread t : threads) { t.join(); } } catch (Exception e) { e.printStackTrace(); } } }
几个通用注意点
- 所有同步方法都会抛出
InterruptedException(CyclicBarrier还会抛出BrokenBarrierException),必须处理或抛出,否则编译不通过。 - 如果某个线程在同步时被中断,可能会导致同步器进入broken状态,需要额外的异常处理来恢复或终止流程。
- 你的
local_counter是实例变量,每个工作线程持有自己的实例,所以不需要额外同步,线程安全没问题。
内容的提问来源于stack exchange,提问作者Kovalainen
相关产品推荐
相关产品推荐

