CyclicBarrier
CyclicBarrier 同步屏障
CyclicBarrier 是可循环使用的屏障,主要功能是让一组线程到达一个屏障时被阻塞,直到最后一个线程到达屏障时,屏障才会打开;所有被屏障拦截的线程才会继续执行。
属性字段
//独占锁 private final ReentrantLock lock = new ReentrantLock(); //等待的条件 private final Condition trip = lock.newCondition(); //线程等待的数量,重置count private final int parties; //被唤醒时,优先执行的任务 private final Runnable barrierCommand; //描述更新换代,重置? //Generation中的broken表示这一次是否完成,默认false。count为0时设置为true private Generation generation = new Generation(); //记录当前需要等待到来的线程数,等于0表示到下一代,通过parties来重置 private int count; |
构造函数
public CyclicBarrier(int parties, Runnable barrierAction) { if (parties <= 0) throw new IllegalArgumentException(); this.parties = parties; this.count = parties; this.barrierCommand = barrierAction; }
public CyclicBarrier(int parties) { this(parties, null); } |
应用场景
CyclicBarrier 适用于多线程结果合并的操作,用于多线程计算数据,最后合并计算结果的应用场景。比如,我们需要统计多个 Excel 中的数据,然后等到一个总结果。我们可以通过多线程处理每一个 Excel ,执行完成后得到相应的结果,最后通过 barrierAction 来计算这些线程的计算结果,得到所有Excel的总和。
使用
|
public class CyclicBarrierTest { public static void main(String[] args) { CyclicBarrier barrier = new CyclicBarrier(3); for (int i = 0; i < 3; i++) { new Writer(barrier).start(); } }
static class Writer extends Thread { private CyclicBarrier cyclicBarrier;
public Writer(CyclicBarrier cyclicBarrier) { this.cyclicBarrier = cyclicBarrier; }
@Override public void run() { System.out.println(Thread.currentThread().getName() + "正在写入数据..."); try { Thread.sleep(2000); //以睡眠来模拟写入数据操作 System.out.println(Thread.currentThread().getName() + "写入数据完毕,等待其他线程写入完毕"); cyclicBarrier.await(); } catch (InterruptedException e) { e.printStackTrace(); } catch (BrokenBarrierException e) { e.printStackTrace(); } } } } Thread-2正在写入数据... Thread-0正在写入数据... Thread-1正在写入数据... Thread-0写入数据完毕,等待其他线程写入完毕 Thread-2写入数据完毕,等待其他线程写入完毕 Thread-1写入数据完毕,等待其他线程写入完毕 |