J.U.C体系进阶(四):juc

Java - J.U.C体系进阶

作者:Kerwin

邮箱:806857264@qq.com

说到做到,就是我的忍道!

juc-sync 同步器框架

同步器名称

作用

CountDownLatch

倒数计数器,构造时设定计数值,当计数值归零后,所有阻塞线程恢复执行;其内部实现了AQS框架

CyclicBarrier

循环栅栏,构造时设定等待线程数,当所有线程都到达栅栏后,栅栏放行;其内部通过ReentrantLock和Condition实现同步

Semaphore

信号量,类似于“令牌”,用于控制共享资源的访问数量;其内部实现了AQS框架

Exchanger

交换器,类似于双向栅栏,用于线程之间的配对和数据交换;其内部根据并发情况有“单槽交换”和“多槽交换”之分

Phaser

多阶段栅栏,相当于CyclicBarrier的升级版,可用于分阶段任务的并发控制执行;其内部比较复杂,支持树形结构,以减少并发带来的竞争

CountDownLatch

注意:CountDownLatch和CyclicBarrier非常相似,且CyclicBarrier是可以重用的,根据具体的场景不同,代码结构不同,其实两者之间可以相互转化,详见CyclicBarrier模块,下文是CountDownLatch-Demo

1// 用法比较简单,直接上代码即可 2// 1.CountDownLatch的同一对象传递 3// 2.构造参数的默认值需要指定 4// 3.线程完成的countDown()->会使默认值减一 5// 4.主线程awiw()等待,所有线程都countDown之后,主线程执行 6// 应用场景:比如五个子线程文件输出导出数据,主线程等所有子线程都完成之后开始压缩操作,上传文件 7 8public class TestCountDownLatch { 9 10 /*** * 关键点:面向对象的方式->参数传递,把CountDownLatch进行传递,使其共用同一个参数 * @param args */ 11 public static void main(String[] args) { 12 CountDownLatch latch = new CountDownLatch(5); 13 14 ExecutorService executorService = Executors.newCachedThreadPool(); 15 MyWoker m1 = new MyWoker("work1", latch); 16 MyWoker m2 = new MyWoker("work2", latch); 17 MyWoker m3 = new MyWoker("work3", latch); 18 MyWoker m4 = new MyWoker("work4", latch); 19 MyWoker m5 = new MyWoker("work5", latch); 20 Boss boss = new Boss("boos", latch); 21 executorService.submit(m1); 22 executorService.submit(m2); 23 executorService.submit(m3); 24 executorService.submit(m4); 25 executorService.submit(m5); 26 executorService.submit(boss); 27 28 executorService.shutdown(); 29 } 30} 31 32class MyWoker implements Callable<String> { 33 34 private String name; 35 private CountDownLatch latch; 36 37 public MyWoker (String name, CountDownLatch latch) { 38 this.name = name; 39 this.latch = latch; 40 } 41 42 @Override 43 public String call() throws Exception { 44 System.out.println(name + " 工人开始工作"); 45 int time = (int)(Math.random() * 100) * 50; 46 Thread.sleep(time); 47 System.out.println(name + " 工人已经完成任务!"); 48 latch.countDown(); 49 return "successful"; 50 } 51 52 public String getName() { 53 return name; 54 } 55 56 public void setName(String name) { 57 this.name = name; 58 } 59 60 public CountDownLatch getLatch() { 61 return latch; 62 } 63 64 public void setLatch(CountDownLatch latch) { 65 this.latch = latch; 66 } 67} 68 69class Boss implements Callable<String> { 70 71 private String name; 72 private CountDownLatch latch; 73 74 public Boss (String name, CountDownLatch latch) { 75 this.name = name; 76 this.latch = latch; 77 } 78 79 @Override 80 public String call() throws Exception { 81 System.out.println("老板准备就绪,等工人都完成了就来视察~"); 82 latch.await(); 83 System.out.println("老板来了,快跑啊~"); 84 return "successful"; 85 } 86 87 public String getName() { 88 return name; 89 } 90 91 public void setName(String name) { 92 this.name = name; 93 } 94 95 public CountDownLatch getLatch() { 96 return latch; 97 } 98 99 public void setLatch(CountDownLatch latch) { 100 this.latch = latch; 101 } 102}

CyclicBarrier

CyclicBarrier是一个同步辅助类,它允许一组线程相互等待直到所有线程都到达一个公共的屏障点

一句话概述就是:“人满发车”

重点理解:

CountDownLatch主要用于主线程阻塞,等待子线程执行完毕后,主线程执行,例如报表导出压缩上传,子线程处理报表,主线程等都执行完毕后,压缩,上传

CyclicBarrier侧重点是人满发车,比如LOL,需要等待是个用户都加载好了之后,再开启主线程执行工作,值得注意的是,这是一般意义的CyclicBarrier

但是,CyclicBarrier提供了另一个构造方法,即可以指定默认额外的执行线程

CyclicBarrier barrier = new CyclicBarrier(5,  new TotalTask(totalService));

这意味着,在很多情况CyclicBarrier可以代替CountDownLatch,主要看代码的结构设计

比如刚才的问题:报表导出压缩上传,子线程处理报表,主线程等都执行完毕后,压缩,上传

如果我用着这种构造方法,配合awit()的位置,让压缩上传线程默认作为最后执行的线程,即可保证执行的顺序,来看个Demo吧:

Demo-1 CyclicBarrier普通使用方法:

1public class TestCyclicBarrier { 2 3 public static void main(String[] args) throws InterruptedException { 4 CyclicBarrier cycli = new CyclicBarrier(10); 5 6 for (int i = 0; i < 9; i++) { 7 new Thread(new BarrierThread("张" + i, cycli)).start(); 8 } 9 10 Thread.sleep(3000); 11 new Thread(new BarrierThread("张" + 10, cycli)).start(); 12 13 Thread.sleep(5000); 14 } 15 16} 17 18class BarrierThread implements Runnable{ 19 20 private String name; 21 private CyclicBarrier cycli; 22 23 public BarrierThread(String name, CyclicBarrier cycli) { 24 super(); 25 this.name = name; 26 this.cycli = cycli; 27 } 28 29 @Override 30 public void run() { 31 System.out.println(name + " 准备就绪"); 32 try { 33 cycli.await(); 34 } catch (InterruptedException | BrokenBarrierException e) { 35 e.printStackTrace(); 36 } 37 System.out.println(name + " 开始执行"); 38 } 39 40 public String getName() { 41 return name; 42 } 43 44 public void setName(String name) { 45 this.name = name; 46 } 47 48 public CyclicBarrier getCycli() { 49 return cycli; 50 } 51 52 public void setCycli(CyclicBarrier cycli) { 53 this.cycli = cycli; 54 } 55}

Demo-2 CyclicBarrier 另一种构造方法的使用:

注意 CyclicBarrier barrier = new CyclicBarrier(5, new TotalTask(totalService));

再注意子线程中,awit等待的代码位置,这个代码位置在程序的最后面,因此CyclicBarrier 的灵活性完全可以由我们来把控,到底在哪一点阻塞,完全是我们自己控制的,这样的变化,就可以在一定程度上替代CountDownLatch,达到更加灵活的目的,CyclicBarrier 且是可重复使用的,细节可以再去深入了解

1/** * 各省数据独立,分库存偖。为了提高计算性能,统计时采用每个省开一个线程先计算单省结果,最后汇总。 * * @author guangbo email:weigbo@163.com * */ 2public class Total { 3 4 // private ConcurrentHashMap result = new ConcurrentHashMap(); 5 6 public static void main(String[] args) { 7 TotalService totalService = new TotalServiceImpl(); 8 CyclicBarrier barrier = new CyclicBarrier(5, 9 new TotalTask(totalService)); 10 11 // 实际系统是查出所有省编码code的列表,然后循环,每个code生成一个线程。 12 new BillTask(new BillServiceImpl(), barrier, "北京").start(); 13 new BillTask(new BillServiceImpl(), barrier, "上海").start(); 14 new BillTask(new BillServiceImpl(), barrier, "广西").start(); 15 new BillTask(new BillServiceImpl(), barrier, "四川").start(); 16 new BillTask(new BillServiceImpl(), barrier, "黑龙江").start(); 17 18 } 19} 20 21/** * 主任务:汇总任务 */ 22class TotalTask implements Runnable { 23 private TotalService totalService; 24 25 TotalTask(TotalService totalService) { 26 this.totalService = totalService; 27 } 28 29 public void run() { 30 // 读取内存中各省的数据汇总,过程略。 31 totalService.count(); 32 System.out.println("======================================="); 33 System.out.println("开始全国汇总"); 34 } 35} 36 37/** * 子任务:计费任务 */ 38class BillTask extends Thread { 39 // 计费服务 40 private BillService billService; 41 private CyclicBarrier barrier; 42 // 代码,按省代码分类,各省数据库独立。 43 private String code; 44 45 BillTask(BillService billService, CyclicBarrier barrier, String code) { 46 this.billService = billService; 47 this.barrier = barrier; 48 this.code = code; 49 } 50 51 public void run() { 52 System.out.println("开始计算--" + code + "省--数据!"); 53 billService.bill(code); 54 // 把bill方法结果存入内存,如ConcurrentHashMap,vector等,代码略 55 System.out.println(code + "省已经计算完成,并通知汇总Service!"); 56 try { 57 // 通知barrier已经完成 58 barrier.await(); 59 } catch (InterruptedException e) { 60 e.printStackTrace(); 61 } catch (BrokenBarrierException e) { 62 e.printStackTrace(); 63 } 64 } 65 66}

相关方法:

— getParties()

获取CyclicBarrier打开屏障的线程数量,也成为方数。

— getNumberWaiting()

获取正在CyclicBarrier上等待的线程数量。

—await()

—await(timeout,TimeUnit)

—isBroken()

获取是否破损标志位broken的值,此值有以下几种情况:

CyclicBarrier初始化时,broken=false,表示屏障未破损。
如果正在等待的线程被中断,则broken=true,表示屏障破损。
如果正在等待的线程超时,则broken=true,表示屏障破损。
如果有线程调用CyclicBarrier.reset()方法,则broken=false,表示屏障回到未破损状态。
—reset()

使得CyclicBarrier回归初始状态,直观来看它做了两件事:

如果有正在等待的线程,则会抛出BrokenBarrierException异常,且这些线程停止等待,继续执行。
将是否破损标志位broken置为false。

CountDownLatch和CyclicBarrier的主要联系和区别如下:

1.闭锁CountDownLatch做减计数,而栅栏CyclicBarrier则是加计数。

2.CountDownLatch是一次性的,CyclicBarrier可以重用。

3.CountDownLatch强调一个线程等多个线程完成某件事情。CyclicBarrier是多个线程互等,等大家都完成。

4.鉴于上面的描述,CyclicBarrier在一些场景中可以替代CountDownLatch实现类似的功能

Semaphore

信号量Semaphore是一个并发工具类,用来控制可同时并发的线程数,其内部维护了一组虚拟许可,通过构造器指定许可的数量,每次线程执行操作时先通过acquire方法获得许可,执行完毕再通过release方法释放许可。如果无可用许可,那么acquire方法将一直阻塞,直到其它线程释放许可

Semaphore用来控制并发线程数,但有个问题FixedThreadPool也可以控制最大并发数,那两者有何不一样呢?首先,量级来看,Semaphore轻量级,是一个并发工具类,线程池重量级无疑,其次特点来讲,Semaphore内的线程是我们实实在在自己创建的,FixedThreadPool是分配给我们的线程池里面的线程

另外,Semaphore如果默认大小为1的时候,还可以当作互斥锁使用,且有公平锁和非公平锁之分(是否按顺序执行,是则就是公平的,但是非常耗性能)

代码Demo,线程队列和Semaphore配合使用

1public class TestQueeThread_2 { 2 static Semaphore semaphore = new Semaphore(1); 3 public static void main(String[] args) throws InterruptedException { 4 System.out.println("begin:" + (System.currentTimeMillis() / 1000)); 5 BlockingQueue<String> myQueue = new ArrayBlockingQueue<String>(1); 6 7 for (int i = 0; i < 100; i++) { 8 new Thread(new Runnable() { 9 @Override 10 public void run() { 11 try { 12 semaphore.acquire(); 13 String log = myQueue.take(); 14 System.out.println(Thread.currentThread().getName() + ":" + doSome(log)); 15 semaphore.release(); 16 } catch (InterruptedException e) { 17 e.printStackTrace(); 18 } 19 } 20 }).start(); 21 } 22 23 for (int i = 0; i < 100; i++) { // 这行代码不能改动 24 String input = i + ""; 25 myQueue.put(input); 26 } 27 } 28 29 public static String doSome(String input) { 30 String output = null; 31 try { 32 Thread.sleep(1000); 33 output = input + ":" + (System.currentTimeMillis() / 1000); 34 return output; 35 } catch (InterruptedException e) { 36 e.printStackTrace(); 37 } 38 return output; 39 } 40} 41 42// 这种用法可能违背了oracle的本意,我们来看看oracle的官方Demo 43// Semaphore是用来约束线程使用共享资源的,控制数据一致性,当然还得是锁 44class Pool { 45 private static final int MAX_AVAILABLE = 100; // 可同时访问资源的最大线程数 46 private final Semaphore available = new Semaphore(MAX_AVAILABLE, true); 47 protected Object[] items = new Object[MAX_AVAILABLE]; //共享资源 48 protected boolean[] used = new boolean[MAX_AVAILABLE]; 49 public Object getItem() throws InterruptedException { 50 available.acquire(); 51 return getNextAvailableItem(); 52 } 53 public void putItem(Object x) { 54 if (markAsUnused(x)) 55 available.release(); 56 } 57 private synchronized Object getNextAvailableItem() { 58 for (int i = 0; i < MAX_AVAILABLE; ++i) { 59 if (!used[i]) { 60 used[i] = true; 61 return items[i]; 62 } 63 } 64 return null; 65 } 66 private synchronized boolean markAsUnused(Object item) { 67 for (int i = 0; i < MAX_AVAILABLE; ++i) { 68 if (item == items[i]) { 69 if (used[i]) { 70 used[i] = false; 71 return true; 72 } else 73 return false; 74 } 75 } 76 return false; 77 } 78}

Exchanger

Exchanger是用作线程并发协作的工具类,简单一句话讲,如果A,B线程都拥有Exchanger对象,如果某一个调用Exchanger的交换方法exchange时候,快的那个会主动等慢的那个(等的意思就是挂起),然后都到位之后,互相唤醒交换数据

代码Demo:

1public class ExchangerTest { 2 static class Producer extends Thread { 3 private Exchanger<Integer> exchanger; 4 private static int data = 0; 5 6 Producer(String name, Exchanger<Integer> exchanger) { 7 super("Producer-" + name); 8 this.exchanger = exchanger; 9 } 10 11 @Override 12 public void run() { 13 for (int i = 1; i < 5; i++) { 14 try { 15 TimeUnit.SECONDS.sleep(1); 16 data = i; 17 System.out.println(getName() + " 交换前:" + data); 18 data = exchanger.exchange(data); 19 System.out.println(getName() + " 交换后:" + data); 20 } catch (InterruptedException e) { 21 e.printStackTrace(); 22 } 23 } 24 } 25 } 26 27 static class Consumer extends Thread { 28 private Exchanger<Integer> exchanger; 29 private static int data = 0; 30 31 Consumer(String name, Exchanger<Integer> exchanger) { 32 super("Consumer-" + name); 33 this.exchanger = exchanger; 34 } 35 36 @Override 37 public void run() { 38 while (true) { 39 data = 0; 40 System.out.println(getName() + " 交换前:" + data); 41 try { 42 TimeUnit.SECONDS.sleep(1); 43 data = exchanger.exchange(data); 44 } catch (InterruptedException e) { 45 e.printStackTrace(); 46 } 47 System.out.println(getName() + " 交换后:" + data); 48 } 49 } 50 } 51 52 public static void main(String[] args) throws InterruptedException { 53 Exchanger<Integer> exchanger = new Exchanger<Integer>(); 54 new Producer("", exchanger).start(); 55 new Consumer("", exchanger).start(); 56 TimeUnit.SECONDS.sleep(7); 57 System.exit(-1); 58 } 59} 60 61结果打印: 62Consumer- 交换前:0 63Producer- 交换前:1 64Consumer- 交换后:1 65Producer- 交换后:0 66Consumer- 交换前:0 67Producer- 交换前:2 68Producer- 交换后:0 69Consumer- 交换后:2 70Consumer- 交换前:0 71Producer- 交换前:3 72Producer- 交换后:0 73Consumer- 交换后:3 74Consumer- 交换前:0 75Producer- 交换前:4 76Producer- 交换后:0 77Consumer- 交换后:4 78Consumer- 交换前:0

我暂时没有碰到用需要此特点的需求,不过很显然它可以用作生产者消费者模式,通过网上的解释来看,Exchanger的实现是非常复杂的,主要是依赖CAS自旋操作

Phase

Phaser 是一个多栅栏的同步工具

phase(阶段) - Phaser也有栅栏,在Phaser中,栅栏的名称叫做phase(阶段),在任意时间点,Phaser只处于某一个phase(阶段),初始阶段为0,最大达到Integerr.MAX_VALUE,然后再次归零。当所有parties参与者都到达后,phase值会递增

parties(参与者) - 其实就是CyclicBarrier中的参与线程的概念,CyclicBarrier中的参与者在初始构造指定后就不能变更,而Phaser既可以在初始构造时指定参与者的数量,也可以中途通过register、bulkRegister、arriveAndDeregister等方法注册/注销参与者

arrive(到达) / advance(进阶) - Phaser注册完parties(参与者)之后,参与者的初始状态是unarrived的,当参与者到达(arrive)当前阶段(phase)后,状态就会变成arrived。当阶段的到达参与者数满足条件后(注册的数量等于到达的数量),阶段就会发生进阶(advance)——也就是phase值+1

1public class SwimmerTest { 2 3 // 游泳选手个数 4 private static int swimmerNum = 6; 5 6 public static void main(String[] args) { 7 8 Phaser phaser = new Phaser(7){ 9 @Override 10 protected boolean onAdvance(int phase, int registeredParties) { 11 System.out.println("---------------PHASE[" + phase + "],Parties[" + registeredParties + "] ---------------"); 12 return phase >= 2 || registeredParties == 0; 13 } 14 }; 15 16 for (int i = 0; i < swimmerNum; i++) { 17 new Thread(new Swimmer(phaser), "swimmer" + i).start(); 18 } 19 20 // 主线程到达,开启第二阶段 21 phaser.arriveAndAwaitAdvance(); 22 23 // 主线程销毁,开启第三阶段 24 phaser.arriveAndDeregister(); 25 26 // 加while是为了防止其它线程没结束就打印了"比赛结束" 27 while (!phaser.isTerminated()) {} 28 29 System.out.println("===== 比赛结束 ====="); 30 } 31} 32 33class Swimmer implements Runnable { 34 35 private Phaser phaser; 36 37 public Swimmer(Phaser phaser) { 38 this.phaser = phaser; 39 } 40 41 @Override 42 public void run() { 43 44 // 从这里到第一个phaser.arriveAndAwaitAdvance()是第一阶段做的事 45 System.out.println("游泳选手-" + Thread.currentThread().getName() + ":已到达赛场"); 46 47 phaser.arriveAndAwaitAdvance(); 48 49 // 从这里到第二个phaser.arriveAndAwaitAdvance()是第二阶段做的事 50 System.out.println("游泳选手-" + Thread.currentThread().getName() + ":已准备好"); 51 52 phaser.arriveAndAwaitAdvance(); 53 54 // 从这里到第三个phaser.arriveAndAwaitAdvance()是第三阶段做的事 55 System.out.println("游泳选手-" + Thread.currentThread().getName() + ":完成比赛"); 56 57 phaser.arriveAndAwaitAdvance(); 58 } 59}

上文说到,Phase阶段概念,且注册数量和到达数量一致的之后,就会进入下一个阶段,代码中即是如此,一开始注册七个指标,游泳子线程会运行到达6个,然后由主线程控制到达,进入到下一个阶段,注意的是参与的线程可以注册也可以销毁,所以主线程阶段二是到达,阶段三测试了销毁

Phaser 的onAdvance方法:

1Phaser phaser = new Phaser(7){ 2 @Override 3 protected boolean onAdvance(int phase, int registeredParties) { 4 System.out.println("---------------PHASE[" + phase + "],Parties[" + registeredParties + "] ---------------"); 5 return phase >= 2 || registeredParties == 0; 6 } 7};

当前阶段,最后一个线程到达后,会触发onAdvance方法,此处是打印了信息,且写明了Phaser终止的标志,注册线程数为0或阶段数到达2 (0,1,2)

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )