预见猿份
主题
首页面试题在线工具关于我们老苗一对一私教学员评价
实战项目
项目前置基础创新WMS项目Java微服务框架与实战云岚到家项目闪聚支付项目学成在线项目青橙电商项目JVM原理与实战调优分布式事务专题Java高频面试题MySQL从入门到精通Java数据结构与算法Java 并发编程(JUC)实战课程老苗一对一私教学员评价blog
blog
  • Java 并发编程(JUC)实战课程

    • 课程介绍
    • 01 多线程入门
    • 02 线程的常用方法
    • 03 线程池
    • 04 synchronized教程
    • 05 ReentrantLock教程
    • 06 读写锁
    • 07 JMM教程
    • 08 CAS教程
    • 09 wait&notify
    • 10 常用的并发集合
    • 11 异步任务控制
    • 12 常用并发工具







----- 到底线了 -----

第12章 常用并发工具 ​

4.1 CountDownLatch ​

4.1.1 CountDownLatch介绍 ​

CountDownLatch是一种倒计时控制多线程的并发工具类,比如下边的例子:线程1需要在线程2和线程3完成工作后才可以继续执行,流程如下:

配图25

1、首先设定计数器(例子: new CountDownLatch(2)表示计数器为2)。

2、线程thread01调用awit()方法等待计算器清零。

3、线程thread02调用CountDownLatch的countDown()方法将计数器减1.

4、线程thread03调用CountDownLatch的countDown()方法将计数器减1.

5、计数器清零唤醒线程thread01继续执行。

CountDownLatch可以实现一个线程等多个线程、多个线程等一个线程,它的应用场景有很多,比如下边的社保查询例子:

配图26

分别用不同的线程查询居民信息和参保信息等信息,最后将每个信息进行汇总组装。

CountDownLatch的Api方法如下:

配图27

4.1.2 CountDownLatch测试 ​

下边模拟两个线程查询参保状态、交费状态,另外一个线程进行信息组装,代码如下:

java
package com.yjoffer.javase.thread.waitnotify;


import java.util.concurrent.CountDownLatch;  
import java.util.concurrent.ExecutorService;  
import java.util.concurrent.Executors;

import java.util.concurrent.TimeUnit;

import javax.xml.bind.helpers.ValidationEventImpl;

/**

- 并发工具测试
- @author 预见猿份(www.yjoffer.com)
- 

*/  
public class ThreadUtilsTest {


//测试CountDownLatch
public static void test_CountDownLatch() {
	//计数器为2
	CountDownLatch countDownLatch = new CountDownLatch(2);
	ExecutorService threadPool = Executors.newCachedThreadPool();
	threadPool.execute(()->{
		logger.debug("组装线程等待中...");
		try {
			countDownLatch.await();
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		logger.debug("组装完成");
	});
	threadPool.execute(()->{
		try {
			TimeUnit.SECONDS.sleep(2);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		logger.debug("参保状态");
		//计数器减1
		countDownLatch.countDown();
	});
	threadPool.execute(()->{
		try {
			TimeUnit.SECONDS.sleep(1);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		logger.debug("交费状态");
		//计数器减1
		countDownLatch.countDown();
	});

}

//休眠
public static void sleep(int time,TimeUnit timeUnit) {
	try {
		timeUnit.sleep(time);
	} catch (InterruptedException e) {
		e.printStackTrace();
	}
}
public static void main(String[] args) {
	test_CountDownLatch1();
}


}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72

输出:

java
21:52:43.491 FINE com.yjoffer.javase.thread2.waitnotify.ThreadUtilsTest[pool-1-thread-1] - 组装线程等待中...
21:52:44.492 FINE com.yjoffer.javase.thread2.waitnotify.ThreadUtilsTest[pool-1-thread-3] - 交费状态
21:52:45.491 FINE com.yjoffer.javase.thread2.waitnotify.ThreadUtilsTest[pool-1-thread-2] - 参保状态
21:52:45.492 FINE com.yjoffer.javase.thread2.waitnotify.ThreadUtilsTest[pool-1-thread-1] - 组装完成
1
2
3
4

4.1.3 CountDownLatch工作原理 ​

查看CountDownLatch类的源代码:

java
    public void await() throws InterruptedException {
        sync.acquireSharedInterruptibly(1);
    }
1
2
3

sync类继承了AbstractQueuedSynchronizer(AQS类),源代码如下:

java
private static final class Sync extends AbstractQueuedSynchronizer {
        private static final long serialVersionUID = 4982264981922014374L;

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

        int getCount() {
            return getState();
        }

        protected int tryAcquireShared(int acquires) {
            return (getState() == 0) ? 1 : -1;
        }

        protected boolean tryReleaseShared(int releases) {
            // Decrement count; signal when transition to zero
            for (;;) {
                int c = getState();
                if (c == 0)
                    return false;
                int nextc = c-1;
                if (compareAndSetState(c, nextc))
                    return nextc == 0;
            }
        }
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27

AQS是ReentrantLock、CountDownLatch等工具类的底层类,我们可以查一下ReentrantLock的源代码,如下:

配图28

AQS字面意思是“抽象队列同步器”,基于AQS自定义同步器可以实现独占锁、共享锁及并发工具。

AQS的工作原理的核心是共享状态state和先进先出双向队列,共享状态state作为多线程访问的共享变量,通过它可以控制锁的获取与释放,FIFO双向队列用于存放等待获取锁的线程,相当于线程的阻塞队列,下边通过一个同步锁的例子说明它的工作过程。如下图:

配图29

1、线程1获取锁,使用CAS将state由0更新为1成功,线程1即Node1作为队列的头结点。

2、线程2获取锁,使用CAS将state由0更新为1失败,线程2即Node2加入队尾等待锁的释放。

3、线程1释放锁,使用CAS将state由1更新为0成功,出队,唤醒它的后继结点Node2,线程2获取锁成功。

4、线程3去获取锁,加入队尾,等待线程2释放锁后被唤醒去获取锁。

CountDownLatch使用的是AQS的共享锁模式。

继续查看await()方法的源代码,最终找到Sync类的获取锁的方法:

java
        protected int tryAcquireShared(int acquires) {
            return (getState() == 0) ? 1 : -1;
        }
1
2
3

state状态为0表示锁空闲可获取锁,否则不允许获取,tryAcquireShared方法返回正数表示获取成功,负数表示失败。

在设定CountDownLatch计数器时指定count,源代码如下:

java
    public CountDownLatch(int count) {
        if (count < 0) throw new IllegalArgumentException("count < 0");
        this.sync = new Sync(count);
    }
1
2
3
4

设定义的count即是state状态的初始值。

只有当state为0才表示计数清零,await线程才可以继续执行,否则进入阻塞队列阻塞。

如何让state清零呢?

其它线程执行countDown()方法表示state减1,源代码如下:

java
    public void countDown() {
        sync.releaseShared(1);
    }

    public final boolean releaseShared(int arg) {
        if (tryReleaseShared(arg)) {
            doReleaseShared();
            return true;
        }
        return false;
    }

    protected boolean tryReleaseShared(int releases) {
            // Decrement count; signal when transition to zero
            for (;;) {
                int c = getState();
                if (c == 0)
                    return false;
                int nextc = c-1;
                if (compareAndSetState(c, nextc))
                    return nextc == 0;
            }
        }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23

在tryReleaseShared方法中找到答案,使用cas让state减1。

如果state等于0则tryReleaseShared方法返回true,唤醒等待的线程。

state等于0再次调用await()方法将不再阻塞。

4.2 CyclicBarrier ​

4.2.1 CyclicBarrier介绍 ​

CyclicBarrier是循环(Cyclic)屏障(Barrier),它可以实现多个线程运行到屏障时互相等待,等所有线程全部到达后再一起继续运行。

它的工作流程如下:

1、初始时设置参与线程数。

2、参与线程到达屏障处等待其它线程。

3、当到达屏障处的线程数达到初始时设定的数量时所有线程继续执行。

4、参与线程到达屏障处等待其它线程。

5、当到达屏障处的线程数达到初始时设定的数量时所有线程继续执行。

6、照这样继续循环...

CyclicBarrier的执行流程如下图:

1、多个线程各自开始执行,每个线程都传入同一个CyclicBarrier对象,参与线程数设置为3。

配图30

2、每个线程执行到屏障处会互相等待,直到全部到达屏障处

所谓屏障处就是每个线程调用awit()的地方。

配图31

3、全部线程到达屏障处继续执行

配图32

4、继续循环一个周期(generation),每个线程执行到屏障处会互相等待,直到全部到达屏障处

配图33

4.2.2 CyclicBarrier测试 ​

CyclicBarrier的api如下:

配图34

配图35

#####

根据CyclicBarrier的运行流程下边进行测试,测试用例为:启动3个线程,到达屏障处再继续运行,照这样循环3次

代码如下:

java
package com.yjoffer.javase.thread.waitnotify;

import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;


/**
 * 并发工具测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadUtilsTest {
	//测试CyclicBarrier 
	public static void test_CyclicBarrier1() {
		//计数器为3
		CyclicBarrier cyclicBarrier = new CyclicBarrier(3, ()->{
			System.out.println("======全部到达屏障=====");
		});
		ExecutorService pool = Executors.newCachedThreadPool();
		pool.execute(()->{
			for (int i = 0; i < 3; i++) {
				try {
					System.out.println("线程1开始执行..");
					TimeUnit.SECONDS.sleep(1);
					System.out.println("线程1到达屏障...");
					cyclicBarrier.await();
				}catch (BrokenBarrierException e) {
					e.printStackTrace();
				}catch (InterruptedException e) {
					e.printStackTrace();
				}
			}

		});
		pool.execute(()->{
			for (int i = 0; i < 3; i++) {
				try {
					System.out.println("线程2开始执行..");
					TimeUnit.SECONDS.sleep(2);
					System.out.println("线程2到达屏障...");
					cyclicBarrier.await();
				}catch (BrokenBarrierException e) {
					e.printStackTrace();
				}catch (InterruptedException e) {
					e.printStackTrace();
				}
			}

		});
		pool.execute(()->{
			for (int i = 0; i < 3; i++) {
				try {
					System.out.println("线程3开始执行..");
					TimeUnit.SECONDS.sleep(3);
					System.out.println("线程3到达屏障...");
					cyclicBarrier.await();
				}catch (BrokenBarrierException e) {
					e.printStackTrace();
				}catch (InterruptedException e) {
					e.printStackTrace();
				}
			}

		});
		pool.shutdown();
	}
	//休眠
	public static void sleep(int time,TimeUnit timeUnit) {
		try {
			timeUnit.sleep(time);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
	}
	public static void main(String[] args) {
		test_CyclicBarrier1();
	}


}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85

经过多次测试结果符合预期,3个线程先后到达屏障后继续执行。

输出:

java
线程1开始执行..
线程2开始执行..
线程3开始执行..
线程1到达屏障...
线程2到达屏障...
线程3到达屏障...
======全部到达屏障=====
线程3开始执行..
线程1开始执行..
线程2开始执行..
线程1到达屏障...
线程2到达屏障...
线程3到达屏障...
======全部到达屏障=====
线程3开始执行..
线程1开始执行..
线程2开始执行..
线程1到达屏障...
线程2到达屏障...
线程3到达屏障...
======全部到达屏障=====
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21

(注意:整体上是先到达屏障才会继续执行,线程1、2、3开始执行与继续执行的次序由于操作系统调度会略有不同)

4.2.3 CyclicBarrier工作原理 ​

1、使用ReentrantLock锁及Condition条件控制对象。

java
public class CyclicBarrier {
...
    /** The lock for guarding barrier entry */
    private final ReentrantLock lock = new ReentrantLock();
    /** Condition to wait on until tripped */
    private final Condition trip = lock.newCondition();
    /** The number of parties */
    private final int parties;
    /* The command to run when tripped */
    private final Runnable barrierCommand;
...
1
2
3
4
5
6
7
8
9
10
11

2、构造一个CyclicBarrier对象,需要指定count及barrierAction。

如下代码

java
    public CyclicBarrier(int parties, Runnable barrierAction) {
        if (parties <= 0) throw new IllegalArgumentException();
        this.parties = parties;
        this.count = parties;
        this.barrierCommand = barrierAction;
    }
1
2
3
4
5
6

barrierAction是当count为0时开始执行。

3、调用await()会将count减1,如果count减1不为0则进入等待队列。

如果count等于0则达到等待线程数量要求,此时唤醒等待队列的所有线程。

4.2.4 秒杀团购案例实现 ​

预见猿份在双十一及毕业季退出秒杀团购活动,规则如下:

1、先到先得,每够10位同学团购成功一次。

2、秒杀活动有一定的时间限制,时间一至活动结束,未达到10位同学则订单作废。

3、采用单独线程池执行提交订单任务 ,因为生产中是向数据库提交订单,而数据库的链接是有限的,所以在高并发的系统中有单独的线程提交订单。

整体结构如下:

配图36

代码如下:

1、学生类

在ThreadUtilsTest内部创建学生类:

java
package com.yjoffer.javase.thread.waitnotify;

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;


/**
 * 并发工具测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadUtilsTest {
	static class PbStudent {

		long id;

		public PbStudent(long id) {
			this.id = id;
		}

		public long getId() {
			return id;
		}

	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25

2、测试方法

java
//测试CyclicBarrier,预见猿份秒杀团购活动
	public static void test_CyclicBarrier2() {

		ExecutorService threadPool = Executors.newCachedThreadPool();

		BlockingQueue<PbStudent> stuQueue = new LinkedBlockingQueue<>();

		CyclicBarrier cyclicBarrier = new CyclicBarrier(10,()->{
			System.out.println("预见猿份团队课程够10位同学成功一次");
		});

		for (int i = 0; i < 98; i++) {
			PbStudent pbStudent = new PbStudent(i);
			threadPool.execute(()->{
				try {
					//等待够10位一组
					cyclicBarrier.await();
					//到达屏障继续执行
					logger.debug("填写报名表");
					stuQueue.put(pbStudent);
				} catch (InterruptedException e) {
					e.printStackTrace();
				} catch (BrokenBarrierException e) {
					e.printStackTrace();
				}
			});
		}
		//订单号
		AtomicLong atomicOrder = new AtomicLong(0);
		//报名成功数量
		AtomicLong count = new AtomicLong(0);
		//由提交订单线程保存订单,为控制并发请求数据库数量将提交 订单线程数量固定
		ExecutorService commitPool = Executors.newFixedThreadPool(20);
		for (int i = 0; i < 20; i++) {
			commitPool.execute(()->{
				while (!Thread.interrupted()){
					try {
						//取出学生
						PbStudent take = stuQueue.take();
						//生成订单号
						long order = atomicOrder.incrementAndGet();
						//保存订单

						TimeUnit.MILLISECONDS.sleep(200);
						logger.debug("报名成功,学号:"+take.getId()+",订单号:"+order);
						count.incrementAndGet();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
			});
		}

		//秒杀10秒
		try {
			TimeUnit.SECONDS.sleep(10);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		System.out.println("团购成功"+count.get());
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61

共模拟请求98人,计数器每达到10打开屏障继续执行,共打开屏障9次,最后8人由于未够10人未打开屏障。

部分输出:

java
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
预见猿份团购课程够10位同学成功一次
...
团购成功90
1
2
3
4
5
6
7
8
9
10
11

4.3 Semaphore ​

4.3.1 Semaphore介绍 ​

Semaphore表示信号量,它可控制并发访问资源。

Semaphore主要应用在并发限制,例如:秒杀限流,针对秒杀活动的高并发量为保证系统稳定可以限制并发线程数,比如限制为50表示最多并发处理50个请求。

Semaphore相当于一个许可证,创建Semaphore时需指定最高可获取许可证的数量,即并发执行的线程数量。

线程先通过acquire方法获取该许可证可以继续执行,执行完成释放许可证,其它线程获取许可证。

Semaphore的常用api如下:

java
//构造方法,permits 表示许可线程的数量
public Semaphore(int permits)		
//表示获取许可证
public void acquire() throws InterruptedException	
//表示释放许可证
public void release()
1
2
3
4
5
6

注意:获取许可证成功,用完一定要释放许可,如下代码:

java
//获取许可
acquire()
try{

}finally{
//释放许可
release()
}
1
2
3
4
5
6
7
8

4.3.2 Semaphore测试 ​

同时开启20个线程,最大获取许可证数量为2。

java
package com.yjoffer.javase.thread.waitnotify;

import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;


/**
 * 并发工具测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadUtilsTest {
	//测试Semaphore
	public static void test_Semaphore1() {
		//许可证数为2
		Semaphore semaphore = new Semaphore(2);
		ExecutorService pool = Executors.newCachedThreadPool();
		for (long i = 0; i <20; i++) {
			pool.execute(()->{
				//申请许可证
				try {
					logger.debug("申请许可证...");
					semaphore.acquire();
					try {
						logger.debug("执行");
						TimeUnit.SECONDS.sleep(3);
						logger.debug("执行完毕");
					}finally {
						//释放许可证
						semaphore.release();
					}
				} catch (InterruptedException e) {
					e.printStackTrace();
				}
			});
		}
	}

	public static void main(String[] args) {
		test_Semaphore1();
	}


}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50

每个线程的处理时间为3秒,运行程序观察输出,当前最多只有两个线程获取许可证并执行。

部分输出:

java
...
pool-1-thread-2执行完毕
pool-1-thread-1执行完毕
pool-1-thread-3执行
pool-1-thread-4执行
pool-1-thread-4执行完毕
pool-1-thread-3执行完毕
pool-1-thread-6执行
pool-1-thread-5执行
pool-1-thread-6执行完毕
pool-1-thread-5执行完毕
pool-1-thread-9执行
pool-1-thread-8执行
pool-1-thread-8执行完毕
pool-1-thread-9执行完毕
...
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16

4.3.3 Semaphore工作原理 ​

Semaphore和CountDownLatch一样底层都是使用的AQS。

1、申请许可证

跟踪申请许可证的方法:

java
    public class Semaphore implements java.io.Serializable {
    private static final long serialVersionUID = -3222578661600680210L;
    /** All mechanics via AbstractQueuedSynchronizer subclass */
    private final Sync sync;

    /**
     * Synchronization implementation for semaphore.  Uses AQS state
     * to represent permits. Subclassed into fair and nonfair
     * versions.
     */
    abstract static class Sync extends AbstractQueuedSynchronizer {
        private static final long serialVersionUID = 1192457210091910933L;

		...
	}
    public void acquire() throws InterruptedException {
        sync.acquireSharedInterruptibly(1);
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18

继续跟踪AQS中的方法,Semaphore使用的是AQS的共享模式:

java
    public final void acquireSharedInterruptibly(int arg)
            throws InterruptedException {
        if (Thread.interrupted())
            throw new InterruptedException();
        //获取许可,如果失败则添加到阻塞队列
        if (tryAcquireShared(arg) < 0)
            doAcquireSharedInterruptibly(arg);
    }
1
2
3
4
5
6
7
8

Semaphore分为公平模式与非公平模式,下边看公平模式:

java
    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;
                //使用cas扣减许可量
                if (remaining < 0 ||
                    compareAndSetState(available, remaining))
                    return remaining;
            }
        }
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23

非公平模式则是不判断阻塞的线程,而是直接进行竞争。

跟踪代码发现:获取许可量就是将state减去一个许可量,减去后的差如果小于0则将线程添加到阻塞队列。

2、释放许可

获取许可量的反方向就是释放许可,如下代码:

java
    public class Semaphore implements java.io.Serializable {
    ...
    public void release() {
        sync.releaseShared(1);
    }
    ...
1
2
3
4
5
6

继续跟踪AQS:

java
public final boolean releaseShared(int arg) {
	//如果释放许可量成功则唤醒下一个结点
    if (tryReleaseShared(arg)) {
        doReleaseShared();
        return true;
    }
    return false;
}
1
2
3
4
5
6
7
8

继续跟踪Semaphore:

java
        protected final boolean tryReleaseShared(int releases) {
            for (;;) {
            	//获取当前许可量
                int current = getState();
                //加上一定的许可量
                int next = current + releases;
                //超过了int的最大值,溢出
                if (next < current) // overflow
                    throw new Error("Maximum permit count exceeded");
                //cas更新增加后的许可量
                if (compareAndSetState(current, next))
                    return true;
            }
        }
1
2
3
4
5
6
7
8
9
10
11
12
13
14

跟踪代码发现:释放许可就是state加一个许可量。

4.3.4 秒杀团购案例优化 ​

下边对预见猿份秒杀团购案例进行优化,控制请求的最大并发数。

(关于秒杀团购案例的内容请参考“cyclicBarrier”章节。)

1、学生类

在ThreadUtilsTest内部创建学生类:

java
package com.yjoffer.javase.thread.waitnotify;

import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;


/**
 * 并发工具测试
 * @author 预见猿份(www.yjoffer.com)
 *
 */
public class ThreadUtilsTest {
	static class PbStudent {

		long id;

		public PbStudent(long id) {
			this.id = id;
		}

		public long getId() {
			return id;
		}

	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25

2、测试方法

在创建线程处理请求前添加获取信号量,达到信号量数量则处理请求。

java
//测试Semaphore,预见猿份秒杀团购活动
	public static void test_Semaphore2() {

		ExecutorService threadPool = Executors.newCachedThreadPool();

		BlockingQueue<PbStudent> stuQueue = new LinkedBlockingQueue<>();

		CyclicBarrier cyclicBarrier = new CyclicBarrier(10,()->{
			System.out.println("预见猿份团队课程够10位同学成功一次");
		});
		//最大并发请求为10
		Semaphore semaphore = new Semaphore(10);
		for (int i = 0; i < 98; i++) {
			PbStudent pbStudent = new PbStudent(i);
			//申请许可证
			semaphore.acquire();
			logger.debug("放行...");
			threadPool.execute(()->{
					try {
							//等待够10位一组
							cyclicBarrier.await();
							//到达屏障继续执行
							logger.debug("填写报名表");
							TimeUnit.SECONDS.sleep(2);
							stuQueue.put(pbStudent);

					} catch (InterruptedException e) {
						e.printStackTrace();
					} catch (BrokenBarrierException e) {
						e.printStackTrace();
					}finally {
						//释放许可
						semaphore.release();
					}
				});

		}
		//订单号
		AtomicLong atomicOrder = new AtomicLong(0);
		//报名成功数量
		AtomicLong count = new AtomicLong(0);
		//由提交订单线程保存订单,为控制并发请求数据库数量将提交 订单线程数量固定
		ExecutorService commitPool = Executors.newFixedThreadPool(20);
		for (int i = 0; i < 20; i++) {
			commitPool.execute(()->{
				while (!Thread.interrupted()){
					try {
						//取出学生
						PbStudent take = stuQueue.take();
						//生成订单号
						long order = atomicOrder.incrementAndGet();
						//保存订单

						TimeUnit.MILLISECONDS.sleep(200);
						logger.debug("报名成功,学号:"+take.getId()+",订单号:"+order);
						count.incrementAndGet();
					} catch (InterruptedException e) {
						e.printStackTrace();
					}
				}
			});
		}

		//秒杀10秒
		try {
			TimeUnit.SECONDS.sleep(40);
		} catch (InterruptedException e) {
			e.printStackTrace();
		}
		System.out.println("团购成功"+count.get());
	}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71

运行代码,跟踪日志发现:一次放行10个请求。

思考:如果Semaphore信号量小于CyclicBarrier循环屏障的参与数量可以吗?

4.4 AQS ​

4.4.1 通过源码分析AQS是什么 ​

AQS是ReentrantLock、CountDownLatch等工具类的底层类,我们可以查一下ReentrantLock的源代码,如下:

配图37

AQS是AbstractQueuedSynchronizer的缩写,字面意思是“抽象队列同步器”,基于AQS自定义同步器可以实现独占锁、共享锁及并发工具,AQS也是企业在面试多线程的必考题目。

AQS的工作原理的核心是共享状态state和FIFO双向队列,共享状态state作为多线程访问的共享变量,通过它可以控制锁的获取与释放,FIFO双向队列用于存放等待获取锁的线程,相当于线程的阻塞队列,下边以ReentrantLock 为例进行说明。

ReentrantLock 是 AQS 独占模式(EXCLUSIVE)最典型的应用:state 的含义是“当前线程重入锁的次数”,state 为 0 表示锁空闲,大于 0 表示已被某个线程持有,其值就是重入次数。ReentrantLock 内部定义的 Sync 类继承 AbstractQueuedSynchronizer,并派生出非公平锁 NonfairSync 与公平锁 FairSync 两个子类,核心源码如下:

java
//ReentrantLock 内部同步器 Sync(简化)
abstract static class Sync extends AbstractQueuedSynchronizer {

    //非公平获取锁:不判断队列中是否有等待线程,直接 CAS 抢锁
    final boolean nonfairTryAcquire(int acquires) {
        final Thread current = Thread.currentThread();
        int c = getState();
        if (c == 0) {                              //锁空闲
            if (compareAndSetState(0, acquires)) { //CAS 将 state 由 0 改为 1
                setExclusiveOwnerThread(current);
                return true;
            }
        } else if (current == getExclusiveOwnerThread()) { //当前线程已持有锁,可重入
            int nextc = c + acquires;              //重入次数加 1
            if (nextc < 0) // overflow
                throw new Error("Maximum lock count exceeded");
            setState(nextc);
            return true;
        }
        return false;
    }

    protected final boolean tryRelease(int releases) {
        int c = getState() - releases;
        if (Thread.currentThread() != getExclusiveOwnerThread())
            throw new IllegalMonitorStateException();
        boolean free = false;
        if (c == 0) {                              //重入次数减到 0 才真正释放锁
            free = true;
            setExclusiveOwnerThread(null);
        }
        setState(c);
        return free;
    }
}

//非公平锁:lock() 时先尝试抢锁,抢不到再排队
static final class NonfairSync extends Sync {
    protected final boolean tryAcquire(int acquires) {
        return nonfairTryAcquire(acquires);
    }
}

//公平锁:队列中已有等待线程时必须排队,保证先来先服务
static final class FairSync extends Sync {
    protected final boolean tryAcquire(int acquires) {
        final Thread current = Thread.currentThread();
        int c = getState();
        if (c == 0) {
            //公平性的关键:hasQueuedPredecessors() 判断队列中是否有前驱结点
            if (!hasQueuedPredecessors() && compareAndSetState(0, acquires)) {
                setExclusiveOwnerThread(current);
                return true;
            }
        } else if (current == getExclusiveOwnerThread()) {
            int nextc = c + acquires;
            if (nextc < 0)
                throw new Error("Maximum lock count exceeded");
            setState(nextc);
            return true;
        }
        return false;
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64

AQS 基于模板方法模式设计,底层把acquire、release、acquireShared、releaseShared 等方式封装成模板方法,把 tryAcquire、tryRelease、tryAcquireShared、tryReleaseShared、isHeldExclusively 等方法留给子类实现。子类只需定义 state 的语义与修改规则,排队与唤醒机制由 AQS 统一完成。

acquire、release、acquireShared、releaseShared的具体含义:

模板方法所属模式执行流程(内部封装了什么)对应工具方法与使用场景
acquire(int arg)独占lock() 最终调用 AQS 的模板方法 acquire(1):先执行子类重写的 tryAcquire 尝试获取资源,获取失败则调用 addWaiter() 把当前线程包装成结点加入 FIFO 队列并阻塞,等前一个线程释放锁后唤醒再继续尝试。
通过 acquireQueued() 自旋:前驱是 head 且再次 tryAcquire 成功则出队拿锁,否则阻塞(park)等待唤醒;
等待期间被中断则补设中断标志。整段流程封装了“抢锁失败 → 排队 → 阻塞 → 唤醒 → 抢锁成功”,子类无需关心。
ReentrantLock.lock() / lockInterruptibly() 的底层调用;对应钩子方法 tryAcquire。
release(int arg)独占1、调用子类重写的 tryRelease(arg) 释放资源;2、释放成功且头结点存在、状态不为 0 时,调用 unparkSuccessor(head) 唤醒后继结点。封装了“释放资源 → 唤醒下一个等待线程”。ReentrantLock.unlock() 的底层调用;对应钩子方法 tryRelease。
acquireShared(int arg)共享1、调用子类重写的 tryAcquireShared(arg):返回值小于 0 表示获取失败;2、失败则通过 doAcquireShared() 入队、自旋、阻塞,直到获取成功;3、获取成功且还有剩余资源时,会向后传播继续唤醒后续线程。CountDownLatch.await()、Semaphore.acquire() 的底层调用;对应钩子方法 tryAcquireShared。
releaseShared(int arg)共享1、调用子类重写的 tryReleaseShared(arg);2、返回 true 时调用 doReleaseShared() 唤醒等待线程,共享模式下会尽可能向后传播唤醒。CountDownLatch.countDown()、Semaphore.release() 的底层调用;对应钩子方法 tryReleaseShared。

acquire、release、acquireShared、releaseShared 是 AQS 对外提供的模板方法(final,不可重写),它们把线程排队、自旋、阻塞、唤醒等通用流程封装好,内部回调子类重写的钩子方法;独占与共享只是同一套 FIFO 队列机制下的两种模式。

补充:除这 4 个基础模板方法外,AQS 还提供了对应的可中断、限时版本(acquireInterruptibly、tryAcquireNanos、acquireSharedInterruptibly、tryAcquireSharedNanos 等),执行流程相同,只是额外增加了中断响应与超时判断。

tryAcquire、tryRelease、tryAcquireShared、tryReleaseShared、isHeldExclusively的具体含义:

钩子方法所属模式返回值与具体含义典型实现(以常用并发工具为例)
tryAcquire(int arg)独占模式返回 true 表示获取资源成功,false 表示失败。arg 表示获取的资源量(一般为 1)。ReentrantLock 中:CAS 把 state 由 0 改为 1,或当前线程已持有锁时 state 加 1(可重入);获取失败时由 AQS 自动把线程包装成结点入队阻塞。
tryRelease(int arg)独占模式返回 true 表示资源已全部释放(锁真正空闲),false 表示还有重入次数未释放完。ReentrantLock 中:state 减 1,减到 0 才返回 true 并清空持有线程,此时 AQS 才会唤醒后继结点。
tryAcquireShared(int arg)共享模式返回负数表示获取失败;返回 0 表示获取成功且无剩余资源;返回正数表示获取成功且有剩余资源,其它线程可以继续获取。CountDownLatch 中:state 为 0 返回 1(计数已清零,await 放行),否则返回 -1(继续阻塞等待);Semaphore 中:剩余许可 = state - 1,小于 0 表示获取失败。
tryReleaseShared(int arg)共享模式返回 true 表示需要唤醒等待线程,false 表示不需要唤醒。CountDownLatch 中:CAS 把 state 减 1,减到 0 时返回 true,唤醒所有 await 的线程;Semaphore 中:state 加 1 后返回 true,唤醒一个等待线程。
isHeldExclusively()独占模式(辅助)返回 true 表示当前线程独占持有资源,false 表示未持有。供 Condition 等需要判断“当前线程是否持有锁”的场景使用,例如 ReentrantLock.newCondition() 的 await / signal 会调用它校验调用者必须持有锁。

其它说明:

1、非公平锁与公平锁的区别只体现在 tryAcquire 上:非公平锁拿到 CPU 时间片就直接 CAS 抢锁,不考虑队列中是否有等待线程;公平锁先调用 hasQueuedPredecessors() 检查队列中是否有前驱结点,有则不抢锁、直接排队,从而实现先来先服务。

2、可重入通过“当前线程 == 持有锁线程”判断:同一线程重复 lock() 时不会入队,只是把 state 加 1;对应的 unlock() 每调用一次 state 减 1,只有减到 0 时才真正释放锁并唤醒后继结点。这也是 ReentrantLock 名字的由来(可重入锁)。

下边以三个线程抢占锁为示例展示三个线程与AQS同步器之间的完整交互过程:

配图38

下图用时序序列图说明详细过程:

wechat_longscreenshot_2026-09-20_123157_615.png

时序图步骤详解

时序图中的每一条箭头都对应 AQS 的一次方法调用或状态变化,与本章末尾“AQS源码分析”中的 acquire、addWaiter、acquireQueued、shouldParkAfterFailedAcquire、release、unparkSuccessor 等方法一一对应。逐步中文含义如下:

步骤时序图中的动作中文含义
1T1 -> AQS:acquire(1)线程1调用 AQS 的模板方法 acquire(1) 申请锁资源,期望通过 CAS 把共享状态 state 从 0 改成 1。
2AQS -> AQS:tryAcquire() CAS state 0 to 1AQS 先调用子类重写的 tryAcquire(),使用 CAS 把 state 由 0 改为 1,成功说明线程1抢到了锁。
3AQS --> T1:获取成功,Node1设为head获取成功,AQS 把线程1对应的结点 Node1 设为 FIFO 队列的头结点(此时队列只有它一个结点,相当于锁的占位)。
4T2 -> AQS:acquire(1)线程2也调用 acquire(1) 申请锁资源。
5AQS -> AQS:tryAcquire() CAS失败tryAcquire() 中 CAS 将 0 改为 1 失败(state 当前为 1,锁被线程1持有),说明线程2抢锁失败。
6AQS -> AQS:addWaiter() 包装Node2加入队尾AQS 调用 addWaiter(),把当前线程2包装成 Node2 结点,通过 CAS 追加到 FIFO 队列队尾。
7AQS -> AQS:acquireQueued() 自旋检查前驱调用 acquireQueued() 进入自旋循环,反复检查自己前驱结点的状态并尝试再次获取锁。
8AQS -> AQS:shouldParkAfterFailedAcquire() 置前驱SIGNAL再次获取失败后,调用 shouldParkAfterFailedAcquire() 把前驱结点状态置为 SIGNAL,表示“前驱释放锁时需要唤醒本结点”。
9AQS --> T2:park() 阻塞等待唤醒AQS 调用 LockSupport.park() 将线程2阻塞挂起,进入等待唤醒状态。
10T1 -> AQS:release(1)线程1执行完临界区代码,调用 release(1) 释放锁资源。
11AQS -> AQS:tryRelease() CAS state 1 to 0AQS 调用子类重写的 tryRelease(),把 state 由 1 改回 0,表示锁已空闲。
12AQS -> AQS:unparkSuccessor(head) 唤醒Node2调用 unparkSuccessor(head) 从头结点开始查找需要唤醒的后继结点,找到 Node2。
13AQS --> T2:unpark() 唤醒线程2AQS 调用 LockSupport.unpark() 唤醒被阻塞的线程2。
14T2 -> AQS:前驱是head且tryAcquire成功,Node2设为head线程2被唤醒后自旋检查:自己的前驱结点是 head 且 tryAcquire() 成功,于是把 Node2 设为新的头结点,线程2成功拿到锁。
15T3 -> AQS:acquire(1)线程3调用 acquire(1) 申请锁资源。
16AQS -> AQS:tryAcquire() 失败(锁被线程2持有)tryAcquire() 获取失败,此时锁刚被线程2持有,state 为 1。
17AQS -> AQS:addWaiter() 包装Node3加入队尾线程3被包装成 Node3 结点加入 FIFO 队列队尾排队。
18AQS --> T3:park() 阻塞,等待线程2释放后再被唤醒线程3被 park() 阻塞挂起,等线程2释放锁后,将按线程2同样的流程被唤醒并尝试获取锁。

通过上述流程可以看到:线程的入队、自旋、阻塞、唤醒等复杂逻辑全部由 AQS 的模板方法自动完成,使用者只需要实现 tryAcquire / tryRelease 等钩子方法,这正是 AQS 的设计精髓。

小结:

AQS是什么?

答:AQS是ReentrantLock、CountDownLatch等工具类的底层类,AQS的工作原理的核心是共享状态state和FIFO双向队列,共享状态state作为多线程访问的共享变量,通过它可以控制锁的获取与释放,FIFO双向队列用于存放等待获取锁的线程。

AQS 基于模板方法模式设计,底层把acquire、release、acquireShared、releaseShared 等方式封装成模板方法,把 tryAcquire、tryRelease、tryAcquireShared、tryReleaseShared、isHeldExclusively 等方法留给子类实现。子类只需定义 state 的语义与修改规则,排队与唤醒机制由 AQS 统一完成

4.4.2 AQS 是什么? ​

答:AQS 全称AbstractQueuedSynchronizer,抽象队列同步器,是ReentrantLock、CountDownLatch、Semaphore、CyclicBarrier等 JDK 并发工具的底层基础类。 AQS 核心由两部分组成:共享状态 state + FIFO 双向阻塞队列。

  • 共享状态state是多线程之间的同步变量,通过 CAS 操作修改 state 的值,以此控制锁的获取与释放;
  • FIFO 双向队列用来存放竞争锁失败、进入阻塞等待的线程。

AQS 基于模板方法模式设计: AQS 父类封装好acquire、release、acquireShared、releaseShared这些模板方法,负责线程入队、阻塞、唤醒、出队等通用排队逻辑; 把tryAcquire、tryRelease、tryAcquireShared、tryReleaseShared、isHeldExclusively这几个钩子方法交给子类去实现。 子类只需要定义 state 变量的含义、以及 state 的修改规则,线程排队、阻塞唤醒这些复杂机制全部由 AQS 统一实现。

补充小要点(面试加分): AQS 分独占模式(ReentrantLock)和共享模式(CountDownLatch/Semaphore)。独占模式同一时刻只允许一个线程获取同步状态;共享模式允许多个线程同时获取同步状态。

4.4.3 AQS的基本使用 ​

AQS同步器作为一个底层类除了理解它的工作原理以外更主要的是学会使用它,下边举例使用AQS实现一个同步锁,Synchronized和ReentrantLock都是同步锁,同步锁即当前只能有一个线程获取锁,没有获取锁的线程则进行等待。

AQS是一个抽象类,需要定义一个子类继承AQS,这个子类通过定义为内部类,在子类中重写AQS的获取锁及释放锁的方法,如下形式:

java
/**
*同步锁
*/
public class SyncLock implements Lock {

    /**
    *同步器 
    */
    class Sync extends AbstractQueuedSynchronizer{
        //尝试获取资源,成功返回true,失败返回false
        @Override
        protected boolean tryAcquire(int arg) {

        }
        //尝试释放资源,成功返回true,失败返回false
        @Override
        protected boolean tryRelease(int arg) {

        }
    }
    @Override
    public void lock() {
       //调用同步器的方法
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
      //调用同步器的方法
    }

    @Override
    public boolean tryLock() {
        return false;
    }

    @Override
    public boolean tryLock(long time, TimeUnit unit) throws InterruptedException {
        return false;
    }

    @Override
    public void unlock() {

    }

    @Override
    public Condition newCondition() {
        return null;
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50

完整的代码如下:

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;

/**
 * 自定义同步锁
 * @author 预见猿份(www.yjoffer.com)
 * @version 1.0
 */
public class SyncLock implements Lock {

    private static Logger logger = Logger.getLogger(SyncLock.class);

    /**
     * 同步器:继承AQS,只需重写获取资源、释放资源两个钩子方法
     */
    class Sync extends AbstractQueuedSynchronizer {

        //尝试获取资源,成功返回true,失败返回false
        @Override
        protected boolean tryAcquire(int arg) {
            //如果使用CAS可以将状态由0改为1则获取资源成功
            if (compareAndSetState(0, 1)) {
                //设置owner为当前线程
                setExclusiveOwnerThread(Thread.currentThread());
                return true;
            }
            return false;
        }

        //尝试释放资源,成功返回true,失败返回false
        @Override
        protected boolean tryRelease(int arg) {
            //释放资源后owner为null
            setExclusiveOwnerThread(null);
            //释放资源后状态为0,此时其它线程才可以获取资源
            setState(0);
            return true;
        }
    }

    //同步器对象
    private Sync sync = new Sync();

    @Override
    public void lock() {
        //调用同步器的acquire方法获取资源,获取失败时AQS会自动将线程入队阻塞
        sync.acquire(1);
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
        sync.acquireInterruptibly(1);
    }

    @Override
    public boolean tryLock() {
        return sync.tryAcquire(1);
    }

    @Override
    public boolean tryLock(long time, TimeUnit unit) throws InterruptedException {
        return sync.tryAcquireNanos(1, unit.toNanos(time));
    }

    @Override
    public void unlock() {
        //调用同步器的release方法释放资源,并唤醒后继结点
        sync.release(1);
    }

    @Override
    public Condition newCondition() {
        return sync.new ConditionObject();
    }

    public static void main(String[] args) {
        SyncLock syncLock = new SyncLock();
        //第一个线程获取锁,执行临界区代码2秒
        new Thread(() -> {
            syncLock.lock();
            try {
                logger.debug("执行临界区代码..");
                TimeUnit.SECONDS.sleep(2);
            } catch (InterruptedException e) {
                e.printStackTrace();
            } finally {
                //释放锁
                logger.debug("释放锁");
                syncLock.unlock();
            }
        }, "t1").start();
        //另一个线程获取锁,需等第一个线程释放后才能获取锁
        try {
            TimeUnit.MILLISECONDS.sleep(200);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        new Thread(() -> {
            syncLock.lock();
            try {
                logger.debug("执行临界区代码...");
            } finally {
                logger.debug("释放锁");
                syncLock.unlock();
            }
        }, "t2").start();
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114

运行输出 :

java
17:07:11.475 FINE com.yjoffer.javase.thread.waitnotify.SyncLock[t1] - 执行临界区代码..
17:07:13.629 FINE com.yjoffer.javase.thread.waitnotify.SyncLock[t1] - 释放锁
17:07:13.629 FINE com.yjoffer.javase.thread.waitnotify.SyncLock[t2] - 执行临界区代码...
17:07:13.629 FINE com.yjoffer.javase.thread.waitnotify.SyncLock[t2] - 释放锁
1
2
3
4

通过测试发现t1和t2不能同时获取锁,t2在2秒后获取锁执行临界区的代码。

在使用AQS的过程中我们并没有去操作FIFO队列元素的入队和出队,这是由AQS去负责,作为使用方只需要重写获取资源、释放资源的方法即可。

4.4.4 AQS实现共享锁 ​

下边使用AQS实现一个共享锁,共享锁即多个线程可以共同拥有锁。实现共享锁则需要使用AQS的tryAcquireShared、tryReleaseShared两个方法。

java
//获取共享资源,返回值:负数表示失败,0表示成功且没有剩余可用资源,正数表示成功且有剩余资源。
int tryAcquireShared(int arg)
// 释放共享资源,释放资源后同时去唤醒其它线程则返回true,否则返回false
boolean tryReleaseShared(int)
1
2
3
4

下边的代码实现了两个线程共同一个锁即同时允许有两个线程拥有锁,思路如下:

1、初始状态值为2。

2、每次获取锁使用CAS将状态减去1,当状态小于0则获取锁失败。

3、每次释放锁则使用CAS将状态加1。

完整代码如下:

java
package com.yjoffer.javase.thread.waitnotify;

import com.yjoffer.javase.config.Logger;

import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.AbstractQueuedSynchronizer;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;

/**
 * 自定义共享锁
 * @author 预见猿份(www.yjoffer.com)
 * @version 1.0
 **/
public class SharedLock implements Lock {

    private static Logger logger = Logger.getLogger(SharedLock.class);


    public SharedLock(int count){
        sharer = new Sharer(count);
    }

    class Sharer extends AbstractQueuedSynchronizer{

        public Sharer(int count){
            setState(count);
        }

        //获取共享锁,
        @Override
        protected int tryAcquireShared(int arg) {
            while (true){
                //取出当前状态
                int state = getState();
                int newState = state - 1;
                if(newState<0) {
                    return newState;
                }
                if(compareAndSetState(state,newState)){
                    return newState;
                }
            }

        }

        @Override
        protected boolean tryReleaseShared(int arg) {
            while (true){
                int state = getState();
                int newState = state + 1;
                if(compareAndSetState(state,newState)){
                    return true;
                }
            }
        }


    }
    private Sharer sharer;
    @Override
    public void lock() {
        sharer.acquireShared(1);
    }

    @Override
    public void lockInterruptibly() throws InterruptedException {
        sharer.acquireSharedInterruptibly(1);
    }

    @Override
    public boolean tryLock() {
        return sharer.tryAcquireShared(1 )>=0;
    }

    @Override
    public boolean tryLock(long time, TimeUnit unit) throws InterruptedException {
        return sharer.tryAcquireSharedNanos(1,unit.toNanos(time ));
    }

    @Override
    public void unlock() {
        sharer.releaseShared(1);
    }

    @Override
    public Condition newCondition() {
        return null;
    }


    public static void main(String[] args) {
        //定义一个共享锁
        SharedLock syncLock = new SharedLock(2);
        //启动10个线程获取共享锁
        for (int i = 0; i < 10; i++) {
            new Thread(()->{
                syncLock.lock();
                try {
                    //执行临界区代码
                    logger.debug("执行临界区代码..");
                    try {
                        TimeUnit.SECONDS.sleep(2);
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }

                } finally {
                    //释放锁
                    logger.debug("释放锁");
                    syncLock.unlock();
                }
            },"t"+i ).start();
        }

    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118

通过测试发现每次有2个线程获取锁并执行临界区的代码,如果在定义共享锁时传入参数为3则实现3个线程可以同时获取锁。

4.4.5 AQS源码分析 ​

下边通过阅读AQS的源码进一步理解它的工作过程,进入AbstractQueuedSynchronizer类主要分析void acquire(int arg)、boolean release(int arg)两个方法。

java
    //获取资源
    public final void acquire(int arg) {
        if (!tryAcquire(arg) &&
            acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
            selfInterrupt();
    }
1
2
3
4
5
6

1、tryAcquire方法实现尝试获取资源

tryacquire方法正是在AQS子类中实现的方法,上边的同步锁例子中在tryacquire方法中使用CAS将状态值由0改为1。if (!tryAcquire(arg) 代码表示如果tryAcquire方法返回true则获取资源成功,此时acquire方法则结束。

2、如果尝试获取资源失败则将线程加入队列

addWaiter方法实现将当前线程包装为Node对象加入队尾。

Node的状态有如下类别:

SIGNAL:当前结点需要唤醒后继线程。

CANCELLED:取消状态,表示不再获取资源。

CONDITION:处于条件队列中,条件达到才能去获取资源。

PROPAGATE:当前结点获取到资源向后传递,在共享模式时使用。

0:初始状态。

addWaiter源代码如下:

java
private Node addWaiter(Node mode) {
    //将当前线程包装成Node
    Node node = new Node(Thread.currentThread(), mode);
    // Try the fast path of enq; backup to full enq on failure
    Node pred = tail;
    if (pred != null) {
        node.prev = pred;
        //使用CAS交换结点添加到队尾 
        if (compareAndSetTail(pred, node)) {
            pred.next = node;
            return node;
        }
    }
    //如果上边没有添加成功使用enq方法添加
    enq(node);
    return node;
}
//通过自旋将结点添加到队尾
private Node enq(final Node node) {
        for (;;) {
            Node t = tail;
            if (t == null) { // Must initialize,队列没有结点先初始化
                if (compareAndSetHead(new Node()))
                    tail = head;
            } else {
                node.prev = t;
                //添加到队尾
                if (compareAndSetTail(t, node)) {
                    t.next = node;
                    return t;
                }
            }
        }
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34

3、加入队尾后通过acquireQueued方法获取资源

acquireQueued的源代码如下:

java
final boolean acquireQueued(final Node node, int arg) {
        boolean failed = true;
        try {
            boolean interrupted = false;
            for (;;) {
            	//取出结点的前驱结点
                final Node p = node.predecessor();
                //如果前驱结点为头结点并且再次获取资源成功
                if (p == head && tryAcquire(arg)) {
                	//将当前结点设置为头结点,原来的头结点出队
                    setHead(node);
                    p.next = null; // help GC
                    failed = false;
                    return interrupted;
                }
                if (shouldParkAfterFailedAcquire(p, node) &&
                    parkAndCheckInterrupt())
                    interrupted = true;
            }
        } finally {
            if (failed)
                cancelAcquire(node);
        }
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24

本方法通过for(;;)实现自旋,if (p == head && tryAcquire(arg))首先判断当前结点的前驱结点是否是head,如果是并且获取资源成功则当前结点变为头结点,原来的头结点出队。

如果if (p == head && tryAcquire(arg))没有成功则执行shouldParkAfterFailedAcquire(p, node)语句,它的作用是判断结点是否应该阻塞,源代码如下:

java
private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
        int ws = pred.waitStatus;
        //如果当前结点状态为SIGNAL则返回true,这表示下一步该结点会被阻塞
        if (ws == Node.SIGNAL)
            /*
             * This node has already set status asking a release
             * to signal it, so it can safely park.
             */
            return true;
        //如果结点状态大于0表示已经放弃获取资源,下边的代码将该结点过虑掉
        if (ws > 0) {
            /*
             * Predecessor was cancelled. Skip over predecessors and
             * indicate retry.
             */
            do {
                node.prev = pred = pred.prev;
            } while (pred.waitStatus > 0);
            pred.next = node;
        } else {
            /*
             * waitStatus must be 0 or PROPAGATE.  Indicate that we
             * need a signal, but don't park yet.  Caller will need to
             * retry to make sure it cannot acquire before parking.
             */
             //使用CAS将状态更新为SIGNAL
            compareAndSetWaitStatus(pred, ws, Node.SIGNAL);
        }
        return false;
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30

当shouldParkAfterFailedAcquire方法返回true则执行parkAndCheckInterrupt方法将阻塞线程,源代码如下:

java
    private final boolean parkAndCheckInterrupt() {
    	//阻塞线程
        LockSupport.park(this);
        return Thread.interrupted();
    }
1
2
3
4
5

头结点获取资源后最后调用release(int arg)去释放资源,源代码如下:

java
public final boolean release(int arg) {
    if (tryRelease(arg)) {
        Node h = head;
        if (h != null && h.waitStatus != 0)
            unparkSuccessor(h);
        return true;
    }
    return false;
}
1
2
3
4
5
6
7
8
9

首先调用tryRelease方法释放资源,tryRelease正是由AQS的子类去实现,上边例子将状态设置为0,这样其它线程才有机会获取资源。

unparkSuccessor方法唤醒后继结点,源代码如下:

java
 private void unparkSuccessor(Node node) {
        /*
         * If status is negative (i.e., possibly needing signal) try
         * to clear in anticipation of signalling.  It is OK if this
         * fails or if status is changed by waiting thread.
         */
        int ws = node.waitStatus;
        if (ws < 0)//将状态设置为0
            compareAndSetWaitStatus(node, ws, 0);

        /*
         * Thread to unpark is held in successor, which is normally
         * just the next node.  But if cancelled or apparently null,
         * traverse backwards from tail to find the actual
         * non-cancelled successor.
         */
        Node s = node.next;
        //如果后继结点为null或结点状态大于0(表示放弃获取资源)
        if (s == null || s.waitStatus > 0) {
            s = null;
            //从尾结点开始搜索结点状态不为null不为正的结点
            for (Node t = tail; t != null && t != node; t = t.prev)
                if (t.waitStatus <= 0)
                    s = t;
        }
        if (s != null)//将找到结点进行唤醒
            LockSupport.unpark(s.thread);
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28

4.4.6 AQS 与常用并发工具对照 ​

AQS 只提供 state 和 FIFO 队列两样基础设施,state 的语义完全由子类定义。下表总结常用并发工具如何基于 AQS 实现(CountDownLatch、Semaphore 的详细分析见 3.4.1、3.4.3):

工具类模式state 语义获取资源钩子方法释放资源钩子方法
ReentrantLock独占0 表示空闲,大于 0 为重入次数tryAcquire:CAS 0→1 或重入 +1tryRelease:减到 0 才释放并唤醒后继
ReentrantReadWriteLock独占+共享高 16 位读锁计数、低 16 位写锁计数tryAcquire / tryAcquireSharedtryRelease / tryReleaseShared
CountDownLatch共享剩余计数,初始为 counttryAcquireShared:state 为 0 返回 1,否则返回 -1tryReleaseShared:CAS 减 1,减到 0 唤醒所有等待线程
Semaphore共享剩余许可数 permitstryAcquireShared:CAS 减 1,小于 0 则失败tryReleaseShared:CAS 加 1

小结:AQS 把排队、阻塞、唤醒这些最复杂且通用的逻辑沉淀在底层,ReentrantLock、CountDownLatch、Semaphore 等工具只需用几行代码定义好 state 的语义,就能获得完整的同步能力——这就是“AQS 是这些工具类底层”的真正含义。

← 11 异步任务控制








如果发现文档内容有错误或排版错乱,请及时联系站长老苗修改,不胜感激。联系我们
关于我们 | 隐私政策 | 豫ICP备2026003386号-4 | 豫公网安备41010202004008号
目录

本页无章节