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

Java多线程之Semaphore、CountDownLatch、CyclicBarrier

金小弟爱编程 2021-05-11
222

Java多线程


  • Java多线程

    • 实例,跑步比赛,多人同时进行

    • 实例,实现多个窗口购买电影票

    • 信号量Semaphore

    • 闭锁CountDownLatch

    • 同步屏幕CyclicBarrier

    • CyclicBarrier 与 CountDownLatch 区别


信号量Semaphore

使用场景:若有m个资源,但有n条线程(n>m),因此同一时刻只能允许m条线程访问资源,此时可以使用Semaphore控制访问该资源的线程数量.


/**
Semaphore是juc中的一个并发工具类
直译信号量,更形象的说法是许可管理证
*/

public class Semaphore implements java.io.Serializable {
private static final long serialVersionUID = -3222578661600680210L;
private final Sync sync;

abstract static class Sync extends AbstractQueuedSynchronizer {
private static final long serialVersionUID = 1192457210091910933L;

Sync(int permits) {
setState(permits);
}

final int getPermits() {
return getState();
}

final int nonfairTryAcquireShared(int acquires) {
for (;;) {
int available = getState();
int remaining = available - acquires;
if (remaining < 0 ||
compareAndSetState(available, remaining))
return remaining;
}
}

protected final boolean tryReleaseShared(int releases) {
for (;;) {
int current = getState();
int next = current + releases;
if (next < current) // overflow
throw new Error("Maximum permit count exceeded");
if (compareAndSetState(current, next))
return true;
}
}

final void reducePermits(int reductions) {
for (;;) {
int current = getState();
int next = current - reductions;
if (next > current) // underflow
throw new Error("Permit count underflow");
if (compareAndSetState(current, next))
return;
}
}

final int drainPermits() {
for (;;) {
int current = getState();
if (current == 0 || compareAndSetState(current, 0))
return current;
}
}
}

/**
* 非公平
*/

static final class NonfairSync extends Sync {
private static final long serialVersionUID = -2694183684443567898L;

NonfairSync(int permits) {
super(permits);
}

protected int tryAcquireShared(int acquires) {
return nonfairTryAcquireShared(acquires);
}
}

/**
* 公平
*/

static final class FairSync extends Sync {
private static final long serialVersionUID = 2014338818796000944L;

FairSync(int permits) {
super(permits);
}

protected int tryAcquireShared(int acquires) {
for (;;) {
if (hasQueuedPredecessors())
return -1;
int available = getState();
int remaining = available - acquires;
if (remaining < 0 ||
compareAndSetState(available, remaining))
return remaining;
}
}
}

//初始化,许可证数量
public Semaphore(int permits) {
sync = new NonfairSync(permits);
}

//初始化,许可证数量、是否公平模式
public Semaphore(int permits, boolean fair) {
sync = fair ? new FairSync(permits) : new NonfairSync(permits);
}

//当前线程尝试去阻塞的获取一个许可证
//此过程是阻塞的,知道出现以下两种情况
//1.获取到了一个去可证,停止等待,继续执行
//2.当前线程被终端,停止等待,继续执行
public void acquire() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}
//当前线程尝试去阻塞的获取n个许可证
public void acquire(int permits) throws InterruptedException {
if (permits < 0) throw new IllegalArgumentException();
sync.acquireSharedInterruptibly(permits);
}

//阻塞获取一个许可证,不允许中断
public void acquireUninterruptibly() {
sync.acquireShared(1);
}

//阻塞获取n个许可证,不允许中断
public void acquireUninterruptibly(int permits) {
if (permits < 0) throw new IllegalArgumentException();
sync.acquireShared(permits);
}

//当前线程尝试去获取一个许可证,非阻塞的
public boolean tryAcquire() {
return sync.nonfairTryAcquireShared(1) >= 0;
}

//当前线程在规定时间内,阻塞的尝试获取一个许可证,超出规定时间,抛出InterruptedException,并停止等待,继续执行
public boolean tryAcquire(long timeout, TimeUnit unit)
throws InterruptedException {
return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout));
}
//当前线程在规定时间内,阻塞的尝试获取n个许可证
public boolean tryAcquire(int permits, long timeout, TimeUnit unit)
throws InterruptedException {
if (permits < 0) throw new IllegalArgumentException();
return sync.tryAcquireSharedNanos(permits, unit.toNanos(timeout));
}

//当前线程释放一个许可证
public void release() {
sync.releaseShared(1);
}

//当前线程释放n个许可证
public void release(int permits) {
if (permits < 0) throw new IllegalArgumentException();
sync.releaseShared(permits);
}

//获取当前可用许可证的数量
public int availablePermits() {
return sync.getPermits();
}

//当前线程获取所有剩余的可用许可证
public int drainPermits() {
return sync.drainPermits();
}

//通过指示减少可用的许可证
protected void reducePermits(int reduction) {
if (reduction < 0) throw new IllegalArgumentException();
sync.reducePermits(reduction);
}

//判断是否公平
public boolean isFair() {
return sync instanceof FairSync;
}

//判断当前Semaphonre对象上是否存在等待许可证的线程
public final boolean hasQueuedThreads() {
return sync.hasQueuedThreads();
}

//获取当前Semaphonre对象上等待许可证的线程的数量
public final int getQueueLength() {
return sync.getQueueLength();
}

//获取当前Semaphonre对象上等待许可证的线程实例
protected Collection<Thread> getQueuedThreads() {
return sync.getQueuedThreads();
}
}


实例,实现多个窗口购买电影票


/**
* 实现电影票购买需求
*/

public class BookService {
//开放3个窗口买票
private static Semaphore semaphore = new Semaphore(3);
//定义总共10张票
private static AtomicInteger number = new AtomicInteger(20);

/**
*
* @param name 购买人姓名
* @param num 票数
* @return
* @throws Exception
*/

public void book(String name,int num) throws Exception {
if (number.get() <= 0){
System.out.println("电影票已经售罄,请改日再来");
throw new Exception("电影票已经售罄,请改日再来");
}
try {
//限制5秒内买不到票则直接返回
semaphore.tryAcquire(5, TimeUnit.SECONDS);
} catch (InterruptedException e) {
e.printStackTrace();
throw new Exception("排队时间太长,请隔日再来");
}
if (number.get() <= 0){
System.out.println("电影票已经售罄,请改日再来");
throw new Exception("电影票已经售罄,请改日再来");
}
synchronized (this){
System.out.println("顾客" + name + "正在窗口"+semaphore.toString()+"进行购买");
if (number.get() <= 0){
System.out.println("电影票已经售罄,请改日再来");
throw new Exception("电影票已经售罄,请改日再来");
}
int temp = num >= number.get() ? number.get() : num;
int temp1 = temp;
while (temp > 0){
number.decrementAndGet();
temp--;
}
System.out.println("顾客:" + name + "本应该购买" + num +"张电影票,实际购买了" + temp1 + "张电影票,剩余"+ number.get() +"张门票");
System.out.println("顾客" + name + "正在窗口"+semaphore.toString()+"购买完成");
}
}

}


闭锁CountDownLatch

主要用于一个线程需要等待其他线程都执行完成之后再执行。
通过计数器进行是实现的,当计数器的值为0时,标识之前的线程已经执行完毕。
典型用法一:

  1. 某一线程在开始运行前等待n个线程执行完毕。将CountDownLatch的计数器初始化为new CountDownLatch(n)

  2. 当线程执行完毕,就将计数器减1 countdownLatch.countDown(),当计数器的值变为0时,在CountDownLatch上await()的线程就会被唤醒
    典型用法二:

  3. 将计数器初始化为1,CountDownLatch(1)

  4. 多个线程都执行countdownlatch.await(),进行等待。

  5. 当主线程countDown(),计数器变为0时,所有线程同时开始工作。


public class CountDownLatch {

//初始化计数器大小
public CountDownLatch(int count) {
if (count < 0) throw new IllegalArgumentException("count < 0");
this.sync = new Sync(count);
}

//线程都等待,等待计数器大小为0时同时执行
public void await() throws InterruptedException {
sync.acquireSharedInterruptibly(1);
}

//等待超时
public boolean await(long timeout, TimeUnit unit)
throws InterruptedException {
return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout));
}

//递减锁存器的计数,如果计数到达零,则释放所有等待的线程。如果当前计数大于零,则将计数减少.
public void countDown() {
sync.releaseShared(1);
}

//获取当前计数器剩余大小
public long getCount() {
return sync.getCount();
}
}


实例,跑步比赛,多人同时进行


/**
* 跑步比赛
*/

public class RunningGameService {
private CountDownLatch runLatch = new CountDownLatch(1);
private CountDownLatch stopLatch = new CountDownLatch(6);
public void running(){

ThreadPoolExecutor executor = new ThreadPoolExecutor(7,
7,
0,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>());
for (int i =0;i< 7;i++){
executor.execute(new Runner());
}
try {
Thread.sleep((long) (Math.random() * 10000));
System.out.println("裁判"+Thread.currentThread().getName()+"即将发布命令");
runLatch.countDown();
System.out.println("裁判"+Thread.currentThread().getName()+"已发送口令,正在等待所有选手到达终点");
stopLatch.await();
System.out.println("所有选手都到达终点");
System.out.println("裁判"+Thread.currentThread().getName()+"汇总成绩排名");
} catch (InterruptedException e) {
e.printStackTrace();
}
}
class Runner implements Runnable{
@Override
public void run() {
try {
System.out.println("运动员"+Thread.currentThread().getName()+ "正在等待裁判发布命令");
runLatch.await();
System.out.println("运动员"+Thread.currentThread().getName()+ "开始跑步");
Thread.sleep(RandomUtils.nextInt(10)*1000);
System.out.println("运动员"+Thread.currentThread().getName()+ "到达终点");
stopLatch.countDown();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
}


同步屏幕CyclicBarrier


可以用于多线程计算数据,最后合并计算结果的场景。
使用场景:若有多条线程,其中一条线程要等其他线程执行完才能执行,那么可以用闭锁;


public class

{


//同步操作锁
private final ReentrantLock lock = new ReentrantLock();
//线程拦截器
private final Condition trip = lock.newCondition();
//每次拦截的线程数
private final int parties;
//换代前执行的线程
private final Runnable barrierCommand;
//标识栅栏的当前代
private Generation generation = new Generation();

//仍在等待的数量,计数器
private int count;

//切换栅栏到下一代
private void nextGeneration() {
// signal completion of last generation
trip.signalAll();
// set up next generation
count = parties;
generation = new Generation();
}

//打翻栅栏
private void breakBarrier() {
generation.broken = true;
count = parties;
trip.signalAll();
}

//阻塞等待的核心方法
private int dowait(boolean timed, long nanos)
throws InterruptedException, BrokenBarrierException,
TimeoutException {
//获取独占锁
final ReentrantLock lock = this.lock;
lock.lock();
try {
//当前代
final Generation g = generation;
//如果当前栅栏损坏了,则抛出异常
if (g.broken)
throw new BrokenBarrierException();
//如果线程中断了,抛出异常
if (Thread.interrupted()) {
//将栅栏设置为True,并通知同一区域的其他的线程
breakBarrier();
throw new InterruptedException();
}

//每次都将计数器减一
int index = --count;
//计数器的值为0之后则需要唤醒所有的线程并转换到下一代(如果有规定runnable需要执行,则需要先执行)
if (index == 0) { // tripped
boolean ranAction = false;
try {
final Runnable command = barrierCommand;
if (command != null)
command.run();
ranAction = true;
//唤醒所有线程转移到下一代
nextGeneration();
return 0;
} finally {
//确保执行失败都会将所有线程唤醒
if (!ranAction)
breakBarrier();
}
}

// loop until tripped, broken, interrupted, or timed out
for (;;) {
try {
//根据参数确实是非超时等待还是超时等待
if (!timed)
trip.await();
else if (nanos > 0L)
nanos = trip.awaitNanos(nanos);
} catch (InterruptedException ie) {
//如果当前线程被打断并且当前栅栏未打断,则打翻栅栏唤醒其他线程
if (g == generation && ! g.broken) {
breakBarrier();
throw ie;
} else {
//若在捕获中断异常前已经完成在栅栏上的等待, 则直接调用中断操作
Thread.currentThread().interrupt();
}
}

//如果线程因为打翻栅栏操作而被唤醒则抛出异常
if (g.broken)
throw new BrokenBarrierException();
//如果线程因为换代操作而被唤醒则返回计数器的值
if (g != generation)
return index;
//如果线程因为时间到了而被唤醒则打翻栅栏并抛出异常
if (timed && nanos <= 0L) {
breakBarrier();
throw new TimeoutException();
}
}
} finally {
//释放锁
lock.unlock();
}
}

//构造器一
// parties:在释放线程之前必须达到的线程数
// barrierAction:最后一个线程执行完成之后继续执行的runnable
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);
}

//返回目标数
public int getParties() {
return parties;
}

//
public int await() throws InterruptedException, BrokenBarrierException {
try {
return dowait(false, 0L);
} catch (TimeoutException toe) {
throw new Error(toe); // cannot happen
}
}

//阻塞等待,直到阻塞等待数等于预期值才全部释放
public int await(long timeout, TimeUnit unit)
throws InterruptedException,
BrokenBarrierException,
TimeoutException {
return dowait(true, unit.toNanos(timeout));
}

//查看是否处于断开状态
public boolean isBroken() {
final ReentrantLock lock = this.lock;
lock.lock();
try {
return generation.broken;
} finally {
lock.unlock();
}
}

//重置
public void reset() {
final ReentrantLock lock = this.lock;
lock.lock();
try {
breakBarrier(); // break the current generation
nextGeneration(); // start a new generation
} finally {
lock.unlock();
}
}

//获取阻塞等待的个数
public int getNumberWaiting() {
final ReentrantLock lock = this.lock;
lock.lock();
try {
return parties - count;
} finally {
lock.unlock();
}
}
}


CyclicBarrier 与 CountDownLatch 区别

CountDownLatchCyclicBarrier
减计数方式加计数方式
计算为0时释放所有等待线程计算是达到预定值时释放所有等待线程
计数为0时,无法重置计数到达预定值时,可以从0重新开始
调用countDown()方法计数减一,调用await()方法只进行阻塞,对计数没任何影响调用await()方法计数加1,若加1后的值不等于构造方法的值,则线程阻塞
不可重复利用可重复利用


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

评论