Phaser类
Phaser类提供了与CyclicBarrier和CountDownLatch差不多的功能。都可以循环设置屏障。
注意:
1、构造方法Phaser(parties) ,其中parties参数,是说明几个线程一组织共同执行。
2、方法arriveAndAwaitAdvance() 在parties数量没有达到时,线程处于阻塞和等待状态。
3、方法arriveAndAwaitAdvance()返回的int类型的值,说明这是第几组线程开始等待。
4、特别注意,线程的个数,应该是parties参数的倍数,否则最后一组线程如果达不到parties的数量时,将呈现为阻塞状态。
1、测试多个线程在arriveAndAwaitAdvance方法中阻塞

以下是测试代码:
package cn.wangjian.concurrent;
import java.util.Random;
import java.util.concurrent.Phaser;
/**
* Phaser与CountDownLatch和CyclicBarrier类似,都是用于创建可重用的<br>
* 障碍。但Phaser更加的灵活<br>
* @author wangjian
*/
public class Demo01_Phaser {
private Phaser phaser;// 1:声明一个Phaser对象
private Random random = new Random();
public Demo01_Phaser() {
// 2:意思是等待三个线程以后再向后执行
// 注意里面的参数:parties是指等待多少个线程共同执行
// 所以,线程的个数,应该是parties的倍数,否则会阻塞永远不会执行
phaser = new Phaser(3);
for (int i = 0; i <9; i++) {
new MyThread().start();
}
}
// 3:声明线程
class MyThread extends Thread {
@Override
public void run() {
System.out.println("线程:【" + getName() + "】开始运行");
long start = System.currentTimeMillis();
sleep1();// 休眠
System.out.println("线程:【" + getName() + "】休眠完成,休眠时间:" + (System.currentTimeMillis() - start));
// 返回的number应该是第几组这个值会不断的进行增加每三个人为一组,从1开始
int num = phaser.arriveAndAwaitAdvance();
System.out.println("第【"+num + "】组人,线程:【" + getName() + "】放行成功,当前时间:" + System.currentTimeMillis());
}
}
public static void main(String[] args) throws Exception {
new Demo01_Phaser();
}
public void sleep1() {
try {
Thread.sleep(1000 * random.nextInt(3));// 任意的休眠最多5秒钟
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
以下某次执行的效果:
线程:【Thread-1】开始运行
线程:【Thread-6】开始运行
线程:【Thread-5】开始运行
线程:【Thread-4】开始运行
线程:【Thread-3】开始运行
线程:【Thread-3】休眠完成,休眠时间:0
线程:【Thread-0】开始运行
线程:【Thread-2】开始运行
线程:【Thread-7】开始运行
线程:【Thread-8】开始运行
线程:【Thread-7】休眠完成,休眠时间:0
线程:【Thread-5】休眠完成,休眠时间:0
第【1】组人,线程:【Thread-5】放行成功,当前时间:1530153332618
线程:【Thread-6】休眠完成,休眠时间:0
第【1】组人,线程:【Thread-3】放行成功,当前时间:1530153332618
第【1】组人,线程:【Thread-7】放行成功,当前时间:1530153332618
线程:【Thread-4】休眠完成,休眠时间:1010
线程:【Thread-8】休眠完成,休眠时间:1010
第【2】组人,线程:【Thread-4】放行成功,当前时间:1530153333628
第【2】组人,线程:【Thread-8】放行成功,当前时间:1530153333628
线程:【Thread-1】休眠完成,休眠时间:1010
线程:【Thread-2】休眠完成,休眠时间:1010
第【2】组人,线程:【Thread-6】放行成功,当前时间:1530153333628
线程:【Thread-0】休眠完成,休眠时间:2009
第【3】组人,线程:【Thread-0】放行成功,当前时间:1530153334627
第【3】组人,线程:【Thread-1】放行成功,当前时间:1530153334627
第【3】组人,线程:【Thread-2】放行成功,当前时间:1530153334627
通过上面的输出可以看出,每三个线程放行一次。因为设置了parties=3。
2、某个线程异常的处理方法
当某个线程在执行过程中,出现异常退出时,因为缺少N个线程执行arriveAndAwaitAdvance方法,最终造成无法达到parties的数量,就会让之前正常在这儿等待的线程一直处于阻塞状态。
同样是上面的代码,模拟某个线程因异常出退出:
package cn.wangjian.concurrent;
import java.util.Random;
import java.util.concurrent.Phaser;
/**
* Phaser与CountDownLatch和CyclicBarrier类似,都是用于创建可重用的<br>
* 障碍。但Phaser更加的灵活<br>
* @author wangjian
*/
public class Demo01_Phaser {
private Phaser phaser;// 1:声明一个Phaser对象
private Random random = new Random();
public Demo01_Phaser() {
// 2:意思是等待三个线程以后再向后执行
// 注意里面的参数:parties是指等待多少个线程共同执行
// 所以,线程的个数,应该是parties的倍数,否则会阻塞永远不会执行
phaser = new Phaser(3);
for (int i = 0; i <3; i++) {
new MyThread().start();
}
}
// 3:声明线程
class MyThread extends Thread {
@Override
public void run() {
System.out.println("线程:【" + getName() + "】开始运行");
long start = System.currentTimeMillis();
sleep1();// 休眠
if(getName().contains("1")) {
throw new RuntimeException("模拟某个线程因异常退出");
}
System.out.println("线程:【" + getName() + "】休眠完成,休眠时间:" + (System.currentTimeMillis() - start));
// 返回的number应该是第几组这个值会不断的进行增加每三个人为一组,从1开始
int num = phaser.arriveAndAwaitAdvance();
System.out.println("第【"+num + "】组人,线程:【" + getName() + "】放行成功,当前时间:" + System.currentTimeMillis());
}
}
public static void main(String[] args) throws Exception {
new Demo01_Phaser();
}
public void sleep1() {
try {
Thread.sleep(1000 * random.nextInt(3));// 任意的休眠最多5秒钟
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
程序执行的结果为:

线程:【Thread-1】开始运行
线程:【Thread-0】开始运行
线程:【Thread-2】开始运行
线程:【Thread-0】休眠完成,休眠时间:0
线程:【Thread-2】休眠完成,休眠时间:0
Exception in thread "Thread-1" java.lang.RuntimeException: 模拟某个线程因异常退出
at cn.wangjian.concurrent.Demo01_Phaser$MyThread.run(Demo01_Phaser.java:29)
通过上面的代码可以看出,某个线程异常退出以后,由于线程的数量达不到parties的值,所以其他两个线程一直会处于阻塞状态。
arriveAndDeregister方法
不过,幸好,Phaser类还有一个arriveAndDeregister()方法,用于取消一个注册parties值。

以下是取消注册的示例:
package cn.wangjian.concurrent;
import java.util.Random;
import java.util.concurrent.Phaser;
/**
* Phaser与CountDownLatch和CyclicBarrier类似,都是用于创建可重用的<br>
* 障碍。但Phaser更加的灵活<br>
*
* @author wangjian
*/
public class Demo01_Phaser {
private Phaser phaser;// 1:声明一个Phaser对象
private Random random = new Random();
public Demo01_Phaser() {
// 2:意思是等待三个线程以后再向后执行
// 注意里面的参数:parties是指等待多少个线程共同执行
// 所以,线程的个数,应该是parties的倍数,否则会阻塞永远不会执行
phaser = new Phaser(3);
for (int i = 0; i < 3; i++) {
new MyThread().start();
}
}
// 3:声明线程
class MyThread extends Thread {
@Override
public void run() {
System.out.println("线程:【" + getName() + "】开始运行");
long start = System.currentTimeMillis();
sleep1();// 休眠
if (getName().contains("1")) {
try {
throw new RuntimeException("模拟某个线程因异常退出");
} catch (Exception e) {
e.printStackTrace(System.out);
System.out.println("取消注册,线程【"+getName()+"】");
int a = phaser.arriveAndDeregister();//取消一个注册的parties
System.out.println("取消:"+a);
throw new RuntimeException(e);//必须要抛出这个异常来停止这个线程,否则还会继续向下执行arriveAndWaitAdvance
}
}
System.out.println("线程:【" + getName() + "】休眠完成,休眠时间:" + (System.currentTimeMillis() - start));
// 返回的number应该是第几组这个值会不断的进行增加每三个人为一组,从1开始
int num = phaser.arriveAndAwaitAdvance();
System.out.println("第【" + num + "】组人,线程:【" + getName() + "】放行成功,当前时间:" + System.currentTimeMillis());
}
}
public static void main(String[] args) throws Exception {
new Demo01_Phaser();
}
public void sleep1() {
try {
Thread.sleep(1000 * random.nextInt(3));// 任意的休眠最多5秒钟
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
运行结果:

线程:【Thread-1】开始运行
线程:【Thread-2】开始运行
线程:【Thread-0】开始运行
java.lang.RuntimeException: 模拟某个线程因异常退出Exception in thread "Thread-1" java.lang.RuntimeException: java.lang.RuntimeException: 模拟某个线程因异常退出
at cn.wangjian.concurrent.Demo01_Phaser$MyThread.run(Demo01_Phaser.java:41)
Caused by: java.lang.RuntimeException: 模拟某个线程因异常退出
at cn.wangjian.concurrent.Demo01_Phaser$MyThread.run(Demo01_Phaser.java:35)
at cn.wangjian.concurrent.Demo01_Phaser$MyThread.run(Demo01_Phaser.java:35)
取消注册,线程【Thread-1】
取消:0
线程:【Thread-0】休眠完成,休眠时间:1015
线程:【Thread-2】休眠完成,休眠时间:2015
第【1】组人,线程:【Thread-2】放行成功,当前时间:1530154986655
第【1】组人,线程:【Thread-0】放行成功,当前时间:1530154986655
通过上面的代码,可以看出,通过arriveAndDeregister()方法,就可以使parties减少,其他等待的线程就可以正常执行。
3、onAdvance方法用于取消屏障
可以通过重写Phaser类的onAdvance方法,取消屏障功能。
phaser = new Phaser(3) {
@Override
protected boolean onAdvance(int phase, int registeredParties) {
System.out.println("onAdvice方法被调");
return true;//返回true为取消屏障,默认为false
}
};
4、注册与取消注册
Phaser类的register方法用于动态的注册一个parties值。而arriveAndDeregister方法用于取消一个注册的值。
还有一个buldRegister方法,用于一次注册多个parties的值。
5、arrive方法
方法arrive的作用是使parties的值增加1,并且不在屏障处等待,直接向下面的代码继继续执行。
小结:
类Phaser提供了动态增减parties计数的功能,这点比CyclicBarrier类操作parties更加方便。
使用Java并发类对线程进行分组控制时,Phaser类比CyclicBarrier类更加强大。




