Condition


Condition

 

Condition是一线程通信工具,表示多线程下参与数争的线程的一,主要负责线境下对线程的挂起和醒工作。

 

方法

// 造成当前线程在接到信号或被中断之前一直于等待状

void await() throws InterruptedException;

// 造成当前线程在接到信号之前一直于等待状。【注意:方法中断不敏感】。

void awaitUninterruptibly();

// 造成当前线程在接到信号、被中断或到指定等待时间之前一直于等待状

//返回表示时间,如果在`nanosTimeout` 之前醒,那返回 `= nanosTimeout - 消耗时间` ,

//如果返回 `<= 0` ,可以

long awaitNanos(long nanosTimeout) throws InterruptedException;

// 造成当前线程在接到信号、被中断或到指定等待时间之前一直于等待状

boolean await(long time, TimeUnit unit) throws InterruptedException;

// 造成当前线程在接到信号、被中断或到指定最后期限之前一直于等待状

//如果没有到指定时间就被通知,返回 true ,否表示到了指定时间,返回返回 false 。

boolean awaitUntil(Date deadline) throws InterruptedException;

// 醒一个等待线程。该线等待方法返回前必须获得与Condition相

void signal();

// 醒所有等待线程。能等待方法返回的线程必须获得与Condition相

void signalAll();

 

原理

 

Condition 内部维护一个条件列,在的情况下,线用 await,线程会被放置在条件列中并被阻。直到用 signal、signalAll 线程,此后线醒,会放入到 AQS 的同步队列,参与争抢锁资源。

AQS是AbstractQueuedSynchronizer的称,翻译过来就是抽象列同器.

 

await

 

用condition.await()方法会使当前线入等待列并,同时线变为等待状。当await()方法返回,一定是得与condition相关联

AQS主要行以下作:

线程1把自己包装成点,waitStatus设为CONDITION(-2),追加到ConditionObject中的条件列(个ConditionObject有一个自己的条件列);

线程1,把state0;

然后醒等待列中head点的下一个点;

AQS  int state=l  exclusiveOwnerThread  Node head  Node tail  ConditionObject  Node firstWaiter  Node head lastWaiter  head  node  node  node  node  waitStatus=CONDITlON  arthinking/ itzhai.com

 

//await

public final void await() throws InterruptedException {

    // 1、线程如果中断,那

    if (Thread.interrupted())

        throw new InterruptedException();

    // 2、将当前线程包装成一个Node点,加入FIFO列中

    Node node = addConditionWaiter();

    // 3、

    int savedState = fullyRelease(node);

    int interruptMode = 0;

    // 4、判断点是否在同步队列(注意非Condition列)中,如果没有,挂起当前线程,因为该线未具

    while (!isOnSyncQueue(node)) {

        // 5、挂起线

        LockSupport.park(this);

        // 6、中断直接返回

        if ((interruptMode = checkInterruptWhileWaiting(node)) != 0)

            break;

    }

    // 7、参与数争(非中断时执行)

    if (acquireQueued(node, savedState) && interruptMode != THROW_IE)

        interruptMode = REINTERRUPT;

 

    // 清理条件列中状态为cancelled的

    if (node.nextWaiter != null) // clean up if cancelled

        unlinkCancelledWaiters();

 

    if (interruptMode != 0)

reportInterruptAfterWait(interruptMode);

    }

signal

 

一个线行了 condition1.signal之后,主要是了以下事情:

把条件列中的第一个点追加到等待列中;

把等待列原来尾点的waitStatusSIGNAL。

然后继续处理自己的事情,自己的事情理完成之后,会醒等待列中head点的下一个线行工作。

 

AQS  int state-I  exclusiveOwnerThread  Node head  Node tail  ConditionObject  Node firstWaiter  Node head lastWaiter  arthinking / itzhai.com  head  node  node  SIGNAL  node  node  node  waitStatus=CONDITION

 

 public final void signal() {

 // 用signal的线程必持有独占

            if (!isHeldExclusively())

                throw new IllegalMonitorStateException();

            Node first = firstWaiter;

            if (first != null)

                doSignal(first);

        }

 

 private void doSignal(Node first) {

            do {

    // 因first上就要被移到同步队列了,所以将first.nextWaiter,作新的firstWatier。

                if ( (firstWaiter = first.nextWaiter) == null)

                    lastWaiter = null;

                    //切断和等待列的关联

                first.nextWaiter = null;

                                // 如果移不成功且有后续节点,那么继续续节点的

            } while (!transferForSignal(first) &&

                     (first = firstWaiter) != null);

        }

 

final boolean transferForSignal(Node node) {

 

    // 判断点是否已在之前被取消了

    if (!compareAndSetWaitStatus(node, Node.CONDITION, 0))

        return false;

 

    // 用 enq 添加到 同步队列的尾部

    Node p = enq(node);

    int ws = p.waitStatus;

    // node 的上一个点 修改 SIGNAL 这样就可以醒自己了

    if (ws > 0 || !compareAndSetWaitStatus(p, ws, Node.SIGNAL))

        LockSupport.unpark(node.thread);

    return true;

}

使用示例

public class LockConditionDemo {

 

    private static ReentrantLock lock = new ReentrantLock();

    private static Condition condition = lock.newCondition();

 

    public static void main(String[] args) throws InterruptedException {

        Thread threadA = new Thread(() -> {

            System.out.println("threadA start");

            lock.lock();

            System.out.println("threadA getLock Running");

            try {

                condition.await();

            } catch (InterruptedException e) {

                e.printStackTrace();

            }

            lock.unlock();

            System.out.println("threadA end");

        });

 

        Thread threadB = new Thread(() -> {

            System.out.println("threadB start");

            lock.lock();

            System.out.println("threadB getLock Running");

            condition.signal();

            lock.unlock();

            System.out.println("threadB end");

        });

 

        threadA.start();

        TimeUnit.SECONDS.sleep(2);

        threadB.start();

    }

}

threadA start

threadA getLock Running

threadB start

threadB getLock Running

threadB end

threadA end

 

 

Await signal

 

package com.iminer.script;

import java.util.concurrent.locks.Condition;

import java.util.concurrent.locks.ReentrantLock;

 

public class AwaitSignal {

    private static ReentrantLock lock = new ReentrantLock();

    private static Condition condition = lock.newCondition();

    private static volatile boolean flag = false;

 

    public static void main(String[] args) {

        Thread waiter = new Thread(new waiter());

        waiter.start();

        Thread signaler = new Thread(new signaler());

        signaler.start();

    }

 

    static class waiter implements Runnable {

        @Override

        public void run() {

            lock.lock();

            try {

                while (!flag) {

                    System.out.println(Thread.currentThread().getName() + "当前条件不足等待");

                    try {

                        condition.await();

                    } catch (InterruptedException e) {

                        e.printStackTrace();

                    }

                }

                System.out.println(Thread.currentThread().getName() + "接到通知条件足");

            } finally {

                lock.unlock();

            }

        }

    }

 

    static class signaler implements Runnable {

        @Override

        public void run() {

            lock.lock();

            try {

                Thread.sleep(5000);

                flag = true;

                condition.signalAll();

            } catch (InterruptedException e) {

                e.printStackTrace();

            } finally {

                lock.unlock();

            }

        }

    }

}

 

 

 

Thread-0当前条件不足等待

Thread-0接到通知条件