Эта статья участвовала в [Пожалуйста, проверьте|У вас есть возможность подать заявку на бесплатные подарки в Nuggets.】Мероприятия.
предисловие
Если вам посчастливилось подать заявку на участие в благотворительных мероприятиях Nuggets, у вас будет возможность получить новую версию значка, предоставленного Nuggets, приняв участие в комментариях.Конкретные подробности лотереи приведены в конце статьи.
Введение в семафор
Semaphore (семафор) — это параллельный класс инструментов в составе пакета JUC, который используется для управления количеством потоков, одновременно обращающихся к критическим ресурсам (общим ресурсам), и гарантирует, что потоки, обращающиеся к критическим ресурсам, могут правильно и разумно использовать общедоступные ресурсы. Как и ReetrantLock, Semaphore реализуется путем прямого или косвенного вызова методов платформы AQS.
Semaphore внутренне поддерживает набор виртуальных лицензий, количество лицензий можно указать через параметры конструктора. Перед доступом к конкретному ресурсу необходимо получить лицензию с помощью метода приобретения.Если количество лицензий равно 0, поток будет заблокирован до тех пор, пока лицензия не будет доступна. После доступа к ресурсу используйте release, чтобы освободить лицензию.
Реализация семафора
Методы, предоставляемые в классе Semaphore:
// 调用该方法后线程会从许可集中尝试获取一个许可
public void acquire()
// 线程调用该方法时会释放已获取的许可
public void release()
// Semaphore构造方法:permits→许可集数量
Semaphore(int permits)
// Semaphore构造方法:permits→许可集数量,fair→公平与非公平
Semaphore(int permits, boolean fair)
// 从信号量中获取许可,该方法不响应中断
void acquireUninterruptibly()
// 返回当前信号量中未被获取的许可数
int availablePermits()
// 获取并返回当前信号量中立即未被获取的所有许可
int drainPermits()
// 返回等待获取许可的所有线程Collection集合
protected Collection<Thread> getQueuedThreads();
// 返回等待获取许可的线程估计数量
int getQueueLength()
// 查询是否有线程正在等待获取当前信号量中的许可
boolean hasQueuedThreads()
// 返回当前信号量的公平类型,如为公平锁返回true,非公平锁为false
boolean isFair()
// 获取当前信号量中一个许可,当没有许可可用时直接返回false不阻塞线程
boolean tryAcquire()
// 在给定时间内获取当前信号量中一个许可,超时还未获取成功则返回false
boolean tryAcquire(long timeout, TimeUnit unit)
Реализация Semaphore представляет собой разделяемую блокировку на основе AQS, которая делится на два режима: честный и несправедливый.
Реализация синхронизации
abstract static class Sync extends AbstractQueuedSynchronizer {
private static final long serialVersionUID = 1192457210091910933L;
//初始的许可数量,同步锁的state保存许可的数量
Sync(int permits) {
setState(permits);
}
//获取许可数量
final int getPermits() {
return getState();
}
//非公平的方式获取共享锁
final int nonfairTryAcquireShared(int acquires) {
for (;;) {
int available = getState();
//可用许可数减去需求的许可数
int remaining = available - acquires;
//许可数大于0时,CAS获取许可
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");
//CAS更新许可,直到成功
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;
}
}
}
Справедливая реализация FairSync
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;
//无前驱节点时,CAS方式获取许可
int available = getState();
int remaining = available - acquires;
if (remaining < 0 ||
compareAndSetState(available, remaining))
return remaining;
}
}
}
Разница между реализацией честной блокировки и недобросовестной блокировки заключается в том, что: для получения блокировки в режиме справедливой блокировки сначала будет вызван метод hasQueuedPredecessors(), чтобы определить, есть ли узел в очереди синхронизации, и если он существует, он напрямую вернет -1 методуAcquireSharedInterruptably(), если (tryAcquireShared(arg)
Недобросовестная реализация синхронизации
static final class NonfairSync extends Sync {
private static final long serialVersionUID = -2694183684443567898L;
NonfairSync(int permits) {
super(permits);
}
//直接使用Sync的nonfairTryAcquireShared()实现,非公平方式获取许可
protected int tryAcquireShared(int acquires) {
return nonfairTryAcquireShared(acquires);
}
}
Конструктор класса NonfairSync с недобросовестной блокировкой в Semaphore завершается на основе вызова конструктора Sync родительского класса, и количество разрешений, переданных при создании объекта Semaphore, в конечном итоге будет передано в состояние идентификатора состояния синхронизации синхронизатора AQS, как следует:
// 父类 - Sync类构造函数
Sync(int permits) {
setState(permits); // 调用AQS内部的set方法
}
// AQS(AbstractQueuedSynchronizer)同步器
public abstract class AbstractQueuedSynchronizer
extends AbstractOwnableSynchronizer {
// 同步状态标识
private volatile int state;
protected final int getState() {
return state;
}
protected final void setState(int newState) {
state = newState;
}
// 对state变量进行CAS操作
protected final boolean compareAndSetState(int expect, int update) {
return unsafe.compareAndSwapInt(this, stateOffset, expect, update);
}
}
Количество разрешений, переданных при создании объекта Semaphore, фактически инициализирует состояние внутри AQS. После инициализации состояние представляет количество доступных лицензий для текущего объекта семафора.
Nonfair Lock NonfairSync Получить лицензию
// Semaphore类 → acquire()方法
public void acquire() throws InterruptedException {
// Sync类继承AQS,此处直接调用AQS内部的acquireSharedInterruptibly()方法
sync.acquireSharedInterruptibly(1);
}
// AbstractQueuedSynchronizer类 → acquireSharedInterruptibly()方法
public final void acquireSharedInterruptibly(int arg)
throws InterruptedException {
// 判断是否出现线程中断信号(标志)
if (Thread.interrupted())
throw new InterruptedException();
// 如果tryAcquireShared(arg)执行结果不小于0,则线程获取同步状态成功
if (tryAcquireShared(arg) < 0)
// 未获取成功加入同步队列阻塞等待
doAcquireSharedInterruptibly(arg);
}
Метод приобретения семафора (acquire()) окончательно завершается вызовом методаAcquireSharedInterruptably() внутри AQS через объект Sync, а функцияAcquireSharedInterruptably() может реагировать на операцию прерывания потока в процессе получения идентификатора состояния синхронизации, если операция не выполняется. fail Если он прерван, сначала вызовите tryAcquireShared(arg), чтобы попытаться получить номер лицензии.Если приобретение прошло успешно, он вернется, чтобы выполнить бизнес, и метод завершится. Если получение не удалось, вызовите doAcquireSharedInterruptably(arg), чтобы добавить текущий поток в очередь синхронизации для блокировки и ожидания. Метод tryAcquireShared(arg) — это метод, предоставляемый AQS, и он не имеет конкретной реализации. Реализация в классе NonfairSync выглядит следующим образом:
// Semaphore类 → NofairSync内部类 → tryAcquireShared()方法
protected int tryAcquireShared(int acquires) {
// 调用了父类Sync中的实现方法
return nonfairTryAcquireShared(acquires);
}
// Syn类 → nonfairTryAcquireShared()方法
abstract static class Sync extends AbstractQueuedSynchronizer {
final int nonfairTryAcquireShared(int acquires) {
// 开启自旋死循环
for (;;) {
int available = getState();
int remaining = available - acquires;
// 判断信号量中可用许可数是否已<0或者CAS执行是否成功
if (remaining < 0 ||
compareAndSetState(available, remaining))
return remaining;
}
}
}
Во-первых, после получения значения состояния отнимите 1, чтобы получить оставшееся значение.Если оно не меньше 0, это означает, что в текущем семафоре еще есть доступные лицензии.Текущий поток начинает пытаться обновить значение состояния с помощью cas , Если cas успешен, это означает, что состояние синхронизации получено успешно, и возвращается к оставшемуся значению. И наоборот, если оставшееся значение меньше 0, это означает, что количество лицензий в семафоре было получено другими потоками, и в настоящее время доступного количества лицензий нет, и непосредственно возвращается оставшееся значение меньше 0. выполнение метода nonfairTryAcquireShared(acquires) завершается и возвращается к методу AcquireSharedInterruptily AQS.(). Когда возвращаемое оставшееся значение меньше 0, выполняется условие if(tryAcquireShared(arg)
// AbstractQueuedSynchronizer类 → doAcquireSharedInterruptibly()方法
private void doAcquireSharedInterruptibly(int arg)
throws InterruptedException {
// 创建节点状态为Node.SHARED共享模式的节点并将其加入同步队列
final Node node = addWaiter(Node.SHARED);
boolean failed = true;
try {
// 开启自旋操作
for (;;) {
final Node p = node.predecessor();
// 判断前驱节点是否为head
if (p == head) {
// 尝试获取同步状态state
int r = tryAcquireShared(arg);
// 如果r不小于0说明获取同步状态成功
if (r >= 0) {
// 将当前线程结点设置为头节点并唤醒后继节点线程
setHeadAndPropagate(node, r);
p.next = null; // 置空方便GC
failed = false;
return;
}
}
// 调整同步队列中node节点的状态并判断是否应该被挂起
// 并判断是否存在中断信号,如果需要中断直接抛出异常结束执行
if (shouldParkAfterFailedAcquire(p, node) &&
parkAndCheckInterrupt())
throw new InterruptedException();
}
} finally {
if (failed)
// 结束该节点线程的请求
cancelAcquire(node);
}
}
// AbstractQueuedSynchronizer类 → setHeadAndPropagate()方法
private void setHeadAndPropagate(Node node, int propagate) {
Node h = head; // 获取同步队列中原本的head头节点
setHead(node); // 将传入的node节点设置为头节点
/*
* propagate=剩余可用许可数,h=旧的head节点
* h==null,(h=head)==null:
* 非空判断的标准写法,避免原本head以及新的头节点node为空
* 如果当前信号量对象中剩余可用许可数大于0或者
* 原本头节点h或者新的头节点node不是结束状态则唤醒后继节点线程
*
* 写两个if的原因在于避免造成不必要的唤醒,因为很有可能唤醒了后续
* 节点的线程之后,还没有线程释放许可/锁,从而导致再次陷入阻塞
*/
if (propagate > 0 || h == null || h.waitStatus < 0 ||
(h = head) == null || h.waitStatus < 0) {
Node s = node.next;
// 避免传入的node为同步队列的唯一节点,
// 因为队列中如果只存在node一个节点,那么后驱节点s必然为空
if (s == null || s.isShared())
doReleaseShared(); // 唤醒后继节点
}
}
лицензия на выпуск
Логика разрешения на освобождение справедливой блокировки согласуется с реализацией несправедливой блокировки, поскольку оба являются подклассами класса Sync, а логика освобождения блокировки заключается в обновлении состояния минус единица и пробуждении потока узла-преемника. . Поэтому конкретная реализация снятия блокировки передается классу Sync.
// Semaphore类 → release()方法
public void release() {
sync.releaseShared(1);
}
// AbstractQueuedSynchronizer类 → releaseShared(arg)方法
public final boolean releaseShared(int arg) {
// 调用子类Semaphore中tryReleaseShared()方法实现
if (tryReleaseShared(arg)) {
doReleaseShared();
return true;
}
return false;
}
Для снятия блокировки вызывается метод Semaphore.release() После вызова этого метода разрешение, удерживаемое потоком, будет снято, а разрешения/состояние увеличатся на единицу.
Как и в предыдущем методе получения лицензии, метод release() Semaphore для освобождения лицензии также выполняется путем косвенного вызова releaseShared(arg) внутри AQS. Поскольку метод releaseShared(arg) AQS является волшебным, окончательная логическая реализация завершается подклассом Sync семафора следующим образом:
// Semaphore类 → Sync子类 → tryReleaseShared()方法
protected final boolean tryReleaseShared(int releases) {
for (;;) {
// 获取AQS中当前同步状态state值
int current = getState();
// 对当前的state值进行增加操作
int next = current + releases;
// 不可能出现,除非传入的releases为负数
if (next < current)
throw new Error("Maximum permit count exceeded");
// CAS更新state值为增加之后的next值
if (compareAndSetState(current, next))
return true;
}
}
По сравнению с логикой получения разрешения это намного проще: достаточно вызвать метод doReleaseShared() после обновления значения состояния, чтобы разбудить следующий поток узла. Но есть два типа потоков, вызывающих метод doReleaseShared():
Одним из них является поток, который освобождает общий номер блокировки/лицензии. Когда метод release() вызывается для освобождения лицензии, его необходимо вызвать, чтобы разбудить следующий поток.
Второй поток — это поток, который только что получил количество общих блокировок/лицензий.
Суммировать
Количество лицензий/счетчиков, переданных при инициализации, в конечном итоге будет косвенно передано в состояние флага состояния синхронизации AQS. Когда поток пытается получить разделяемую блокировку, состояние уменьшается на 1. Когда состояние равно 0, это означает, что доступной разделяемой блокировки нет, а другие потоки последующих запросов будут инкапсулированы как разделяемые узлы и добавлены к очередь синхронизации для ожидания удержания других блокировок Освобождение потока (состояние плюс один). Отличие от монопольного режима заключается в том, что в совместно используемом режиме, в дополнение к потоку, который пробуждает узел-преемник при снятии блокировки, поток, который успешно получает разделяемую блокировку, также пробуждает узел-преемник при определенных условиях. Подробнее об AQS см. в предыдущей статье:Блокировка Java: подробное объяснение AQS (1) Блокировка Java: подробное объяснение AQS (2)
Инструкции по розыгрышу
1. Участвуйте в комментариях (обсуждайте технологический контент);
2. Правила лотереи: если комментарии в этой статье соответствуют требованиям активности Nuggets, каждый из двух лучших игроков в области горячих комментариев получит новую версию значка Nuggets (если нет горячих комментариев, два счастливчика получат быть выбранным из области комментариев);