暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Java并发控制之Phaser-重用屏障

Coding On Road 2018-06-28
311

Phaser

Phaser类提供了与CyclicBarrierCountDownLatch差不多的功能。都可以循环设置屏障。

注意:

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减少,其他等待的线程就可以正常执行。

 

3onAdvance方法用于取消屏障

可以通过重写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的值。

 

5arrive方法

方法arrive的作用是使parties的值增加1,并且不在屏障处等待,直接向下面的代码继继续执行。

 

 

小结:

 Phaser提供了动态增减parties计数的功能,这点比CyclicBarrier类操作parties更加方便。

使用Java并发类对线程进行分组控制时,Phaser类比CyclicBarrier类更加强大。

 

 


文章转载自Coding On Road,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论