Непонятный AQS (Часть 1)

Java

В блоге мы рассмотрели исходный код CopyOnWriteArrayList, это не сложно, в нем используется реентерабельная эксклюзивная блокировка: ReentrantLock, которая похожа на Synchronized, но полностью отличается от Synchronized.

Характеристики синхронизированных блокировок гарантируются JVM, а характеристики блокировок ReentrantLock контролируются кодом Java верхнего уровня. Основой ReentrantLock является AQS.На самом деле многие параллельные контейнеры используют ReentrantLock, который косвенно использует AQS, и параллельные платформы, такие как CountDownLatch, CyclicBarrier и Semaphore, также используют AQS, что показывает важность AQS.

Однако разобраться в AQS немного глубже непросто, и это включает в себя множество вещей, поэтому этот блог будет разделен на две части.В первой части будут представлены необходимые знания об AQS: LockSupport, основные концепции AQS, и в монопольном и совместно используемом режимах анализ основного исходного кода AQS и т. д., во второй части будет представлена ​​поддержка AQS для условных переменных, а также применение AQS.

Чтобы узнать больше об AQS, вы должны сначала освоить необходимое условие: LockSupport.

LockSupport

LockSupport — это инструментальный класс. Его основная функция — приостанавливать и пробуждать потоки. Его нижний слой — это нативный метод, который вызывается. Мы не будем углубляться в это, а в основном рассмотрим применение LockSupport.

припарковаться, разпарковаться

Если поток, вызвавший метод park, получил лицензию, связанную с LockSupport, он вернется сразу после вызова park, в противном случае поток будет заблокирован до получения лицензии. Если поток вызывает метод unpark, он получает лицензию, связанную с LockSupport. Если поток ранее вызывал парк и был заблокирован, он будет разбужен. Если поток ранее не вызывал метод park, после вызова метода park , он вернется немедленно.

    public static void main(String[] args) {
        System.out.println("Hello,LockSupport");
        LockSupport.park();
        System.out.println("Bye,LockSupport");
    }

результат операции:image.pngПоток печатает первое предложение и блокируется, поскольку поток не получил лицензию, связанную с LockSupport.

    public static void main(String[] args) {
        System.out.println("Hello,LockSupport");
        LockSupport.unpark(Thread.currentThread());
        LockSupport.park();
        System.out.println("Bye,LockSupport");
    }

результат операции:image.pngСначала вызовите метод unpark, передайте в текущем потоке, что текущий поток получил лицензию, связанную с LockSupport, а затем вызовите метод park, поскольку у потока уже есть лицензия, поэтому он немедленно возвращается и выводит второе предложение.

    public static void main(String[] args) {
        Thread thread=new Thread(()->{
            System.out.println("Hello,LockSupport");
            LockSupport.park();
            System.out.println("Bye,LockSupport");
        });
        thread.start();
        LockSupport.unpark(thread);
    }

результат операции:image.pngСначала создается Thread, внутри вызывается метод park, затем запускается поток, в основном потоке вызывается метод unpark и передается дочерний поток. Есть два случая для этого метода:

  • Основной поток сначала вызывает метод unpark, а метод park в дочернем потоке вызывает его позже.
  • Метод парковки дочернего потока вызывается первым, а метод unpark основного потока вызывается позже.

Но в любом случае конечный результат одинаков. Просто процесс другой, первый случай После того, как основной поток вызывает метод unpark, дочерний поток получает лицензию, и дочерний поток возвращается сразу после внутреннего вызова park.Заблокировано, затем основной поток вызывает unpark, позволяет дочернему потоку получить лицензию, и дочерний поток возвращается .

parkNanos(long nanos)

Аналогично методу park, разница в том, что есть дополнительный таймаут, при вызове parkNanos поток блокируется, после превышения nanos он будет возвращен независимо от того, получено разрешение или нет.

    public static void main(String[] args) {
        System.out.println("Hello,LockSupport");
        LockSupport.parkNanos(Integer.MAX_VALUE);
        System.out.println("Bye,LockSupport");
    }

результат операции:image.pngЧтобы увидеть очевидный эффект, я установил время в Integer.MAX_VALUE, видно, что хотя метод unpark не вызывался для получения лицензии, метод возвращался через определенный промежуток времени.

park(Object blocker)

Этот способ рекомендуется, так как при нем информацию о блокирующих объектах можно посмотреть через команду jstack.

public class Main {
    public void test() {
        LockSupport.park(this);
    }
    public static void main(String[] args) {
        Main main = new Main();
        main.test();
    }
}

Используйте команду jstack pid:image.png

Есть несколько методов, которые не будут представлены один за другим.

С вышеупомянутой основой мы можем, наконец, перейти к сегодняшней теме: AQS.

Что такое АКСС

Полное название AQS — AbstractQueuedSynchronizer, что в переводе означает очередь абстрактной синхронизации на китайском языке. Когда я впервые столкнулся с AQS, я впервые почувствовал, что эта штука как-то связана с абстракцией из-за Abstract. . . Позже я обнаружил, что эта вещь не имеет ничего общего с абстракцией.Потихоньку приходит новое понимание.Эта вещь действительно связана с абстракцией, потому что она абстрагирует некоторые методы реализации синхронных очередей для других верхних уровней.Переписывание или повторное использование компонентов. Дело в том, что другим компонентам верхнего уровня нужно переписать свои методы! Более подробно, другие компоненты должны наследовать AbstractQueuedSynchronizer и переписать некоторые методы.

Давайте сначала посмотрим на UML-диаграмму AQS:image.png

Основные концепции AQS

Давайте сначала дадим общее введение в AQS и разберемся в основных вещах в AQS.

AQS поддерживает двунаправленную очередь FIFO Что такое FIFO? Это означает "первым пришел, первым вышел". Двусторонняя очередь означает, что когда предыдущий узел указывает на следующий узел, следующий узел также указывает на предыдущий узел. Мы можем видеть это в классе Node, связанном с AbstractQueuedSynchronizer: prev сохраняет текущий узел Предыдущий узел и следующий хранят следующий узел текущего узла.Существует профессиональный термин, который является узлом-предшественником и узлом-преемником.В то же время класс AbstractQueuedSynchronizer имеет два поля, одно из которых является головным, а другое это хвост Как следует из названия, голова сохраняется Головной узел, хвост сохраняет хвостовой узел.

SHARED в классе Node используется, чтобы отметить, что поток помещается в очередь ожидания, когда он получает общие ресурсы, а EXCLUSIVE используется, чтобы отметить, что поток помещается в очередь ожидания, когда он получает эксклюзивные ресурсы.Из этого предложения мы можно увидеть. Класс Node фактически сохраняет потоки, помещенные в очередь ожидания, а некоторые потоки помещаются в очередь ожидания из-за невозможности получения общих ресурсов, а некоторые потоки помещаются в очередь ожидания из-за невозможности получения эксклюзивных ресурсов. , так тут надо марку различать.

Другими словами, двунаправленная очередь FIFO на самом деле является очередью ожидания в AQS.

В классе Node также есть поле: waitStatus, которое имеет пять значений, а именно:

  • СИГНАЛ: значение равно -1.После того, как текущий узел войдет в очередь и перед переходом в состояние сна, обязательно измените его тип предыдущего узла на СИГНАЛ, чтобы текущий узел можно было разбудить, когда последний отменяется или освобождается.
  • ОТМЕНА: значение равно 1, если оно отменено, поток, ожидающий в очереди ожидания, истекает по времени или прерывается, и узел, входящий в это состояние, больше не будет меняться.
  • УСЛОВИЕ: значение равно -2, узел находится в очереди условий. Когда другие потоки вызывают метод signal() условия, узел переводится в очередь ожидания AQS. Следует отметить, что очередь условий и ожидание очереди AQS не связаны друг с другом.Не одно и то же.
  • РАСПРОСТРАНЕНИЕ: Значение равно -3. Что касается того, что делает это состояние, то в большинстве блогов в Интернете, включая книги, просто упоминается, что оно посвящено режиму обмена и связано с общением, но более глубокого объяснения нет. Беспомощный, я до сих пор не мог понять значение этого значения статуса.
  • 0: значение по умолчанию.

В классе AbstractQueuedSynchronizer есть поле состояния, которое для обеспечения видимости помечено как volatile.Дизайн этого поля потрясающий. Для ReentrantLock состояние сохраняет количество повторных входов Для ReentrantReadWriteLock состояние сохраняет количество повторных входов для получения блокировки чтения и количество повторных входов для блокировки записи.

В классе AbstractQueuedSynchronizer также есть внутренний класс: ConditionObject, который используется для обеспечения поддержки условных переменных.

AQS предоставляет два способа получения ресурсов: один — монопольный, а другой — совместно используемый.

Как упоминалось выше, вам нужно определить класс, наследующий класс AbstractQueuedSynchronizer и переопределяющий методы в нем.

  • Для эксклюзивного режима вам необходимо переопределить методы tryAcquire(arg) и tryRelease(int arg).
  • Для общего режима вам необходимо переопределить методы tryAcquireShared(arg) и tryReleaseShared(int arg).

Анализ исходного кода

эксклюзивный режим

acquire
    public final void acquire(int arg) {
        if (!tryAcquire(arg) &&
            acquireQueued(addWaiter(Node.EXCLUSIVE), arg))
            selfInterrupt();
    }

Этот метод является методом верхнего уровня для получения ресурсов в эксклюзивном режиме. Если поток успешно вызывает метод tryAcquire(arg), это означает, что ресурс был получен и возвращается напрямую. В случае неудачи текущий поток инкапсулируется в a Node вставляется с waitStaus как Node.EXCLUSIVE в конец очереди ожидания AQS.

Давайте посмотрим на метод tryAcquire:

    protected boolean tryAcquire(int arg) {
        throw new UnsupportedOperationException();
    }

Нани, я прямо сообщил об ошибке, что это за хрень? Не забывайте, нам нужно переопределить этот метод.

Давайте посмотрим на метод addWaiter(Node.EXCLUSIVE), arg):

    private Node addWaiter(Node mode) {
        Node node = new Node(Thread.currentThread(), mode);//封装成Node,新的Node
        // Try the fast path of enq; backup to full enq on failure
        Node pred = tail;//把尾节点赋值给pred ,pred也就是尾节点了
        if (pred != null) {//如果pred不为NULL
            node.prev = pred;//pred赋值给新节点的前驱节点,也就是新节点的前驱节点是尾节点
            if (compareAndSetTail(pred, node)) {//CAS,如果pred还是尾节点,则把新节点设置成尾节点,设置成功后,进入if
                pred.next = node;//把新节点赋值给pred的后继节点
                return node;//返回新节点
            }
        }
        enq(node);
        return node;
    }

Этот метод сначала инкапсулирует поток в (Node.EXCLUSIVE Node, сначала попробуйте поместить этот Node непосредственно в конец очереди, в случае успеха вернитесь напрямую, в случае неудачи вызовите enq(node) для входа в операцию очереди:

    private Node enq(final Node node) {
        for (;;) {//自旋
            Node t = tail;//把尾节点赋值给t
            //如果尾节点为空,则新建一个空的Node,用CAS把空的Node设置成头节点
            //成功后,再把尾部节点也指向空的Node
            if (t == null) { // Must initialize
                if (compareAndSetHead(new Node()))
                    tail = head;
            } else {
                node.prev = t;//把尾节点赋值给传进来的node的前驱节点
                if (compareAndSetTail(t, node)) {//CAS,如果t还是尾部节点,则用传进来的node替换旧的尾部节点
                    t.next = node;//设置t的后继节点为传进来的node
                    return t;
                }
            }
        }
    }

Как правило, этот метод заключается в том, чтобы поместить узел, которому не удалось получить ресурсы, в очередь ожидания AQS.

Вернемся к методу верхнего уровня и посмотрим на методAcquireQueued:

    final boolean acquireQueued(final Node node, int arg) {
        boolean failed = true;
        try {
            boolean interrupted = false;
            for (;;) {
                final Node p = node.predecessor();//拿到node的前驱节点,赋值给p
                if (p == head && tryAcquire(arg)) {//如果p已经是头节点了,代表这个时候
//node是第二个节点,再次调用tryAcquire获取资源
                    setHead(node);//设置头节点
                    p.next = null; // help GC
                    failed = false;
                    return interrupted;
                }
                if (shouldParkAfterFailedAcquire(p, node) &&//判断此node是否可以被park
                    parkAndCheckInterrupt())//park
                    interrupted = true;
            }
        } finally {
            if (failed)
                cancelAcquire(node);
        }
    }

Это снова спин CAS. Сначала получите узел-предшественник узла и назначьте его p. Если p уже является головным узлом, это означает, что узел является вторым узлом в это время. Попробуйте еще раз вызвать tryAcquire для получения ресурсов. Если успешный, установить головной узел в узел, вернуть бит флага прерывания, в случае неудачи сначала определить, можно ли его запарковать, если да, запарковать и дождаться разпарковки.

Давайте взглянем на метод parkAndCheckInterrupt:

    private static boolean shouldParkAfterFailedAcquire(Node pred, Node node) {
        int ws = pred.waitStatus;//拿到前驱节点的waitStatus,赋值给ws
        if (ws == Node.SIGNAL)//如果是SIGNAL
            return true;
        if (ws > 0) {//如果是ws>0,则说明前驱节点被取消了,通过while循环,
            //找到最近的一个没有取消的节点,排到后面
            do {
                node.prev = pred = pred.prev;
            } while (pred.waitStatus > 0);
            pred.next = node;
        } else {
            compareAndSetWaitStatus(pred, ws, Node.SIGNAL);//CAS设置前驱节点的waitStatus为SIGNAL
        }
        return false;
    }

Если состояние ожидания узла-предшественника равно SIGNAL, возвращайтесь напрямую. Если узел-предшественник отменен, через цикл while найдите ближайший узел, который не был отменен, и поместите его на задний план. Если узел-предшественник находится в других состояниях , передайте CAS узлу-предшественнику. Состояние ожидания узла установлено в SIGNAL.

Давайте взглянем на метод parkAndCheckInterrupt:

    private final boolean parkAndCheckInterrupt() {
        LockSupport.park(this);
        return Thread.interrupted();
    }

Этот метод относительно прост, то есть сам припарковывается и возвращает значение, если текущий поток прерван.

Подведем итоги для методаAcquireQueued: Найдите безопасную точку, чтобы припарковаться самостоятельно, если она проснется, проверьте, являетесь ли вы второй нодой, если да, попробуйте снова получить ресурс, в случае успеха установите себя головной нодой.

Что ж, весь основной контент верхнего уровня проанализирован, давайте подведем итоги:

  1. Старайтесь быстро получить ресурсы, в случае успеха возвращайтесь сразу, при неудаче переходите к следующему шагу.
  2. Выполните операцию постановки в очередь.
  3. Поток в очереди ожидания получает ресурс.

Наконец, нарисуйте блок-схему, чтобы понять весь процесс:image.png

release
    public final boolean release(int arg) {
        if (tryRelease(arg)) {
            Node h = head;
            if (h != null && h.waitStatus != 0)
                unparkSuccessor(h);
            return true;
        }
        return false;
    }

Этот метод является методом верхнего уровня для освобождения ресурсов в монопольном режиме. Сначала вызовите метод tryRelease. В случае успеха назначьте головной узел h. Если h не равно null и waitStatus не равен 0, вызовите метод unparkSuccessor, чтобы разбудить следующий узел.

попробуйтеВыпуск:

    protected boolean tryRelease(int arg) {
        throw new UnsupportedOperationException();
    }

Этот метод по-прежнему напрямую сообщает об ошибке, потому что нам нужно его переписать. Здесь нам нужно обратить особое внимание, этот метод должен определить, был ли полностью освобожден ресурс, если блокировка повторно используемая, блокировка могла быть получена несколько раз, поэтому последняя блокировка должна быть освобождена до возврата true, иначе он возвращает ложный .

unparkSuccessor:

    private void unparkSuccessor(Node node) {
        int ws = node.waitStatus;//拿到当前节点的waitStatus,赋值给ws
        if (ws < 0)
            compareAndSetWaitStatus(node, ws, 0);
        Node s = node.next;//当前节点的下一个节点赋值给s
        if (s == null || s.waitStatus > 0) {//如果s==null或者已经被取消了,就通过for循环找到下一个需要被唤醒的节点
            s = 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);//唤醒
    }

Суть этого метода заключается в том, чтобы разбудить следующий узел, который необходимо разбудить.

общий режим

acquireShared
      public final void acquireShared(int arg) {
        if (tryAcquireShared(arg) < 0)
            doAcquireShared(arg);
    }

Этот метод является методом верхнего уровня для получения ресурсов в совместно используемом режиме. Сначала вызовите tryAcquireShared, чтобы попытаться получить ресурс.Если это не удается, вызовите doAcquireShared, чтобы войти в очередь ожидания, пока ресурс не будет получен.

попробуйте получить общий доступ:

    protected boolean tryAcquire(int arg) {
        throw new UnsupportedOperationException();
    }

Нам нужно переопределить метод tryAcquireShared.

сделатьAcquireShared:

   private void doAcquireShared(int arg) {
       final Node node = addWaiter(Node.SHARED);//入队
       boolean failed = true;
       try {
           boolean interrupted = false;
           for (;;) {
               final Node p = node.predecessor();//拿到当前节点的前驱节点,赋值给p
               if (p == head) {//如果p是头节点
                   int r = tryAcquireShared(arg);//调用tryAcquireShared尝试获取资源
                   if (r >= 0) {
                       setHeadAndPropagate(node, r);//设置头节点,如果还有剩余资源,唤醒下一个节点
                       p.next = null; // help GC
                       if (interrupted)
                           selfInterrupt();
                       failed = false;
                       return;
                   }
               }
               if (shouldParkAfterFailedAcquire(p, node) &&
                   parkAndCheckInterrupt())
                   interrupted = true;
           }
       } finally {
           if (failed)
               cancelAcquire(node);
       }
   }

Этот метод мало чем отличается от процесса в монопольном режиме, самое большое отличие — это метод setHeadAndPropagate, посмотрим, что делает этот метод:

    private void setHeadAndPropagate(Node node, int propagate) {
        Node h = head; 
        setHead(node);//设置头节点
        //如果还有剩余资源
        if (propagate > 0 || h == null || h.waitStatus < 0 ||
            (h = head) == null || h.waitStatus < 0) {
            Node s = node.next;//找到后继节点
            if (s == null || s.isShared())
                doReleaseShared();//调用doReleaseShared方法
        }
    }

Сначала установите текущий узел в качестве головного узла. Если есть оставшиеся ресурсы, найдите узел-преемник и вызовите метод doReleaseShared. Мы рассмотрим этот метод позже, но из имени метода мы можем знать, что он связан с освобождением общего доступа. Ресурсы.

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

Этот метод является методом верхнего уровня для высвобождения ресурсов в совместно используемом режиме. Метод tryReleaseShared еще нужно переписать, в случае успеха вызываем метод doReleaseShared:

    private void doReleaseShared() {
        for (;;) {
            Node h = head;//把头节点赋值给h
            if (h != null && h != tail) {
                int ws = h.waitStatus;//拿到h的waitStatus赋值给ws
                if (ws == Node.SIGNAL) {//如果为SIGNAL
                    if (!compareAndSetWaitStatus(h, Node.SIGNAL, 0))
                        continue;         
                    unparkSuccessor(h);//唤醒后继节点
                }
                else if (ws == 0 &&
                         !compareAndSetWaitStatus(h, 0, Node.PROPAGATE))
                    continue;            
            }
            if (h == head) 
                break;
        }
    }

Этот метод также называется setHeadAndPropagate в doAcquireShared в методе высшего уровня для получения общих ресурсов в совместно используемом режиме.

Что ж, был проанализирован основной процесс получения и высвобождения ресурсов в монопольном и совместно используемом режимах.

Если вы внимательны, то должны обнаружить, что в AQS есть два метода:AcquireInterruptably()/acquireSharedInterruptably(). Эти два метода являются еще одним Interruptably с точки зрения имен. Они будут реагировать на прерывания, и мы представили их выше. игнорирует прерывания.

Этот блог закончился, но кое-что еще не упомянуто: поддержка условных переменных, эта часть контента будет подробно представлена ​​в следующем блоге.