源码层面解读 ObjectMonitor

1. Mark Word

​ Java 对象在堆内存中的存储布局可以分为三部分:对象头、实例数据、对齐填充。其中,对象头包含了对象在运行时所需的一些元数据,它主要由两部分组成:

​ Mark Word:记录对象自身的运行时数据。

​ 类型指针:它指向元空间中该对象对应的类元数据(Klass 结构体),JVM 通过这个指针来判断该对象是哪个类的实例。

​ 在这里我们需要重点关注的是 Mark Word,理解它是理解 synchronized、锁状态升级等概念的关键。Mark Word 是一个动态的存储空间,根据对象状态的不同存储不同的信息,这主要分为(对象状态括号后的信息表示:是否偏向锁 锁标志位):

​ 无锁状态(0 01):对象的 hashcode + GC 分代年龄

​ 偏向锁状态(1 01):ThreadID + Epoch + GC 分代年龄

​ 轻量级锁状态(0 00):指向线程栈中锁记录的指针

​ 重量级锁状态(0 10):指向 Moniter 的指针

​ GC 标记状态(0 11):无

2. Moniter


并不是每个 Java 对象都关联了一个 Moniter。从 Mark Word 存储信息的角度可以看出,只有该对象进入重量级锁状态的时候,JVM 才会为其分配一个 Monitor 对象,并将对象头中的 Mark Word 指向这个 Monitor 对象,这是一个按需分配的过程。因此,正确的说法是,Moniter 是锁在重量级状态下的具体实现,如果对象还未进入这个状态,是不会关联任何 Moniter 的。
在 HotSpot JVM 中,Moniter 是通过一个名为 ObjectMonitor 的 C++ 数据结构来实现的,源码如下:

// hotspot/src/share/vm/runtime/objectMonitor.hpp
class ObjectMonitor {
    
    // 指向当前持有该 Monitor 的线程,此时该线程为 RUNNABLE 状态。如果为 null 则表示锁被释放
    void * volatile _owner;
    
    // 前一个持有者的线程 ID
    volatile jlong _previous_owner_tid;
    
    /*
     * 锁的重入次数。synchronized 是可重入锁,线程每次进入同一个锁保护的代码块时,值会递增,释放时递减。
     * 只有当值减到 0 时,线程才会真正释放锁。
     */
    volatile intptr_t _recursions;

    // 由这两个队列维护处于锁竞争但尚未获取到锁的线程(BLOCKED 状态的线程)。
    // 当锁被释放时,OS 以某种策略(根据 QMode 的值)从这两个队列中选取一个线程来唤醒。
    ObjectWaiter * volatile _cxq; // 后进先出
    ObjectWaiter * volatile _EntryList; // 先进先出

    // 维护所有在这个对象上调用了 wait() 方法的线程(WAITING 或 TIMED_WAITING 状态的线程)。
    // 当其他线程调用 notify() 时,线程被移至 _cxq 或 _EntryList 中参与锁竞争
    ObjectWaiter * volatile _WaitSet;

    // 其他字段 ...
};	

​ 在 JDK 的早期版本中,所有等待锁的线程都放在 EntryList 中,由于其队列的特性,可以保证一定的公平性,但吞吐量较差。因为先加入锁等待队列的线程更可能保有 CPU 缓存,如果让它直接入队可能不会在短时间内被唤醒。因此为了性能优化引入了 cxq,它是一个栈结构,保存最近的锁等待线程。EntryList 中的线程均由 cxq 以某种方式转移而来,这取决于 QMode 的配置。具体释放锁时的 JDK 源码会在下面展示。

3. Java 重量级锁

​ Moniter 是 Java 实现重量级锁的方式,那么它到底是如何运作的?它是怎样与底层 OS 的 Mutex 和 Condition 交互的?这需要从 JVM 源码找出答案。

​ 当线程尝试进入 synchronized 保护的代码块时,JVM 会执行 ObjectMonitor::enter(TRAPS),当线程退出代码块时,JVM 会执行 ObjectMonitor::exit(bool not_suspended, TRAPS),下面就来看看这两个方法的源码(以 HotSpot 为例,简化源码中的部分非主要流程)。

​ 源码地址:jdk8u/jdk8u/hotspot: 782f3b88b5ba src/share/vm/runtime/objectMonitor.cpp

void ATTR ObjectMonitor::enter(TRAPS) {
    Thread * const Self = THREAD ;
    void * cur ;
    
    // CAS 尝试直接获取锁,如果成功返回 NULL,如果失败返回当前持有锁的线程
    // 这里使用了三个操作数 (新值,指针,旧值),
    // 意思是当且仅当 &_owner 为 NULL 的时候用 Self 更新 &_owner
    cur = Atomic::cmpxchg_ptr (Self, &_owner, NULL) ;
    
    // CAS 成功,锁空闲,直接获取成功
    if (cur == NULL) {   
       return ;
    }
    
    // CAS 失败,但当前持有锁的就是自己,这是锁重入
    if (cur == Self) {
       _recursions ++ ;
       return ;
    }
    
    // 检查当前线程是否持有对应的轻量级锁,轻量级锁持有者无需重新竞争
    if (Self->is_lock_owned ((address)cur)) {
      assert (_recursions == 0, "internal state error");
      _recursions = 1 ;
      _owner = Self ;
      OwnerIsThread = 1 ;
      return ;
    }
    
    // 尝试自旋获取锁
    if (Knob_SpinEarly && TrySpin (Self) > 0) {
       assert (_owner == Self, "invariant") ;
       assert (_recursions == 0, "invariant") ;
       Self->_Stalled = 0 ;
       return ;
    }
    
    // 以上方式都没有获取到锁,这意味着我们遇到了严重的竞争!
    
    // 增加竞争线程计数
    Atomic::inc_ptr(&_count);
    
    for (;;) {
        jt->set_suspend_equivalent();
        
        // 锁竞争的核心逻辑
        EnterI(THREAD);

        // 如果线程在 EnterI 内部被挂起过,这里不会退出循环,而是处理挂起状态
        // 正常情况下竞争到锁后循环退出
        if (!ExitSuspendEquivalent(jt)) break;

        _recursions = 0;
        _succ = NULL;
        exit(false, Self); // 释放刚获得的锁
        jt->java_suspend_self(); // 挂起自己
        // 当线程恢复后,循环继续
    }
    
    // 线程已经成功获取锁,不再需要等待该监视器,清空该字段
    Self->set_current_pending_monitor(NULL);
}

​ 可以看出,即使锁对象关联了一个 Moniter,依然不会立刻竞争 Mutex 挂起线程,而是尝试 CAS、自旋等方式获取锁。并且我们发现,实际上自旋锁的代码是存在于 ObjectMonitor::enter(TRAPS) 的,因此自旋的逻辑是不属于轻量级锁状态的,而是属于重量级锁状态。

​ 并且,新加入竞争的线程,直接调用 Enter 方法可能会先于等待队列中的线程获取到锁。因为在 Enter 方法中,线程会先尝试通过 CAS 操作获取锁,如果获取失败,才会进入 EnterI 方法竞争 Mutex,这意味着新来的线程有可能在等待队列中的线程被唤醒之前抢到锁,这就是 synchronized 非公平的体现。

​ 如果遇到严重的竞争,会进入 EnterI 方法:

void ObjectMonitor::EnterI(TRAPS) {
    Thread * const Self = THREAD;
    assert(Self->is_Java_thread(), "invariant");
    assert(((JavaThread *) Self)->thread_state() == _thread_blocked, "invariant");

    // 再次尝试 CAS 获取锁
    if (TryLock (Self) > 0) {
        assert(_succ != Self, "invariant");
        assert(_owner == Self, "invariant");
        assert(_Responsible != Self, "invariant");
        return;
    }

    // 再次尝试自旋获取锁
    if (TrySpin (Self) > 0) {
        assert(_owner == Self, "invariant");
        assert(_succ != Self, "invariant");
        assert(_Responsible != Self, "invariant");
        return;
    }
    
    // 依然没有获取到锁,准备将线程入队

    // 将当前线程封装成 ObjectWaiter 节点
    ObjectWaiter node(Self);
    Self->_ParkEvent->reset();
    node._prev   = (ObjectWaiter *) 0xBAD;
    node.TState  = ObjectWaiter::TS_CXQ;

    // 使用 CAS 将节点插入到 _cxq 的头部
    ObjectWaiter * nxt;
    for (;;) {
        node._next = nxt = _cxq;
        if (Atomic::cmpxchg_ptr (&node, &_cxq, nxt) == nxt) {
            break;
        }

        // 如果插入失败,再次尝试 CAS 获取锁
        if (TryLock (Self) > 0) {
            assert(_succ != Self, "invariant");
            assert(_owner == Self, "invariant");
            assert(_Responsible != Self, "invariant");
            return;
        }
    }

    // 如果这是第一个竞争线程,或该线程有继任者(_succ),
    // 在当前没有 Responsible 线程的情况下,设置自己为 Responsible 线程。
    if ((nxt == NULL && _EntryList == NULL) || 
        _succ != NULL || 
        _Responsible == Self) {
        if (_Responsible == NULL) {
            Atomic::cmpxchg_ptr (Self, &_Responsible, NULL);
        }
    }

    int nWakeups = 0;
    for (;;) {
        // 再次尝试 CAS 获取锁
        if (TryLock(Self) > 0) {
            break;
        }

        // 依然未获取到锁,执行到这意味着线程真正要被阻塞在 OS 内核中了,也就是 park()
        if (_Responsible == Self) {
            /* 
             * 如果系统没有明确的继任者,意味着此时拿到锁的线程不保证一定会明确唤醒某个线程,
             * 可能有死锁风险,因此 Responsible 线程会使用超时等待。
             * 如果 Responsible 线程超时醒来,它会再次尝试获取锁,如果获取不到,
             * 它会重新选举_Responsible。
             */
            if (_succ == NULL)
                // 带超时的阻塞,线程被真正阻塞在 OS 内核中
                Self->_ParkEvent->park((jlong)1000); // 1秒
            }
        } else {
            // 无限期阻塞,线程被真正阻塞在 OS 内核中
            Self->_ParkEvent->park();
        }

        // 被唤醒后,尝试获取锁
        if (TryLock(Self) > 0) {
            break;
        }

        // 重新选举 Responsible 线程
        if (_Responsible == Self) {
            if (_succ != NULL) {
                _Responsible = NULL;
            }
        }

        // 如果当前线程被中断,则中断等待
        if (Self->is_interrupted(false)) {
            // 从队列中移除节点
            UnlinkAfterAcquire(Self, &node);
            if (Self->is_interrupted(true)) {
                return;
            }
            break;
        }
    }

    // 获取锁成功,从队列中移除节点
    UnlinkAfterAcquire(Self, &node);
    assert(_owner == Self, "invariant");
    return;
}

​ 可以看出,在 EnterI 方法中,线程接触到了 JVM 内部用于线程挂起的更底层的机制 Self->_ParkEvent->park();,关于 ParkEvent 具体是怎么做的,会在后面进行讲解,在此之前我们先来看看线程执行完同步块中的代码后是如何退出的(在这里我只截取了部分从队列唤醒相关的代码)。

void ATTR ObjectMonitor::exit(bool not_suspended, TRAPS) {
    
    for (;;) {
        // QMode = 2 的情况
        if (QMode == 2 && _cxq != NULL) {
            w = _cxq;  // 直接从 cxq 唤醒
            ExitEpilog(Self, w);
            return;
        }
        
        // QMode = 3 或 4 的情况,在这里省略
        if (QMode == 3 && _cxq != NULL) { /* ... */ }
        if (QMode == 4 && _cxq != NULL) { /* ... */ }
        
        // 除去上面的情况,都优先检查 EntryList,如果其不为空就从头部唤醒
        w = _EntryList;
        if (w != NULL) {
            ExitEpilog(Self, w);
            return;
        }
        
        // 若 EntryList 为空,就要从 cxq 转移节点,这又分为两种情况
        w = _cxq;
        if (w == NULL) continue;
        
        for (;;) {
            ObjectWaiter * u = (ObjectWaiter *) Atomic::cmpxchg_ptr(NULL, &_cxq, w);
            if (u == w) break;
            w = u;
        }
        
        if (QMode == 1) { // QMode == 1 时,反转 cxq 顺序再转移
            ObjectWaiter * s = NULL ;
            ObjectWaiter * t = w ;
            ObjectWaiter * u = NULL ;
            while (t != NULL) {
                guarantee (t->TState == ObjectWaiter::TS_CXQ, "invariant") ;
                t->TState = ObjectWaiter::TS_ENTER ;
                u = t->_next ;
                t->_prev = u ;
                t->_next = s ;
                s = t;
                t = u ;
            }
             _EntryList  = s ;
             assert (s != NULL, "invariant") ;
        } else { // QMode == 0 或 QMode == 2,直接转移
            _EntryList = w;
            ObjectWaiter * q = NULL;
            ObjectWaiter * p;
            for (p = w; p != NULL; p = p->_next) {
                p->TState = ObjectWaiter::TS_ENTER;
                p->_prev = q;
                q = p;
            }
        }
        
        w = _EntryList;
        if (w != NULL) {
            // 现在已经选出了将要唤醒的线程,调用这个方法进行 unpark()
            ExitEpilog(Self, w);
            return;
        }
    }
}

// 这个方法完成了 Thread * Self 的释放锁,并唤醒 ObjectWaiter * Wakee
void ObjectMonitor::ExitEpilog (Thread * Self, ObjectWaiter * Wakee) {
    assert (_owner == Self, "invariant") ;

    // Knob_SuccEnabled 是一个 JVM 调优参数,控制是否启用继承者优化。
    // 如果启用优化,将唤醒的线程设为继承者。
    _succ = Knob_SuccEnabled ? Wakee->_thread : NULL ;
    
    // 从将要唤醒的线程中提取对应的线程挂起事件
    ParkEvent * Trigger = Wakee->_event ;
    Wakee  = NULL ;

    // 将 _owner 置为 NULL,意味着当前 Moniter 没有线程占用了,也就是释放了锁
    // 这里插入了全内存屏障,确保锁释放操作对所有处理器可见,并防止指令重排序
    OrderAccess::release_store_ptr (&_owner, NULL) ;
    OrderAccess::fence() ;
    
    // 通过 ParkEvent 唤醒指定线程
    Trigger->unpark() ;

    //...
}

​ 从上面的源码中,我们可以看到,JVM 会从 EntryList 或 cxq 中选择一个线程,然后调用 unpark 来唤醒这个特定的线程。这意味着,唤醒选择是由 JVM 层面显式控制的。

​ 在默认 QMode = 0 时,释放锁的线程会从 EntryList 头部取出线程唤醒,如果 EntryList 中没有线程,cxq 中的线程会直接转移至 EntryList。由于 cxq 本身是非公平的(栈),因此这又一次证明 synchronized 的非公平性。

​ 看到这里相信大家都会有疑惑,因为我们并没有看到任何与 OS Mutex 和 Condition 相关的代码,线程究竟是怎么被挂起在 OS 内核中的?park 方法到底在背后做了什么?接下来我们就继续深究 JVM 的 ParkEvent。

4. ParkEvent

​ 当我们在代码中使用到 synchronizedwait()notify(),它们在底层最终都委托给 ParkEvent 来执行线程的阻塞与唤醒。

ParkEvent 属于 HotSpot JVM 中与操作系统平台相关的底层同步原语实现。不同的 OS 提供了不同的线程挂起/唤醒原语(在 Linux 中称为 Mutex 和 Condition),因此不同的 OS 有不同的 PlatformEvent 实现,ParkEvent 通过继承它们来为 JVM 上层提供统一的 park()unpark() 接口。

​ 在前面的讲解中,我们了解到 Monitor 中有三个队列用于管理处于等待状态的队列,即 EntryList、cxq、WaitSet。每个在这些队列中等待的线程,都会关联一个自己的 ParkEvent 对象,所以你可以把 ParkEvent 看做每个线程自带的挂起/唤醒工具。

​ 这就解释了为什么 Monitor 可以精确唤醒一个线程,因为每个 ParkEvent 对象只负责一个线程的阻塞和唤醒,它并不管理队列,它只是提供了一个简单的挂起和唤醒机制,而其上层的 Monitor 就可以利用 ParkEvent 提供的能力,来使用 EntryList 和 cxq 来管理这些线程的唤醒逻辑。

​ 下面我们就来看一下 ParkEvent 的 C++ 数据结构。

class ParkEvent: public os::PlatformEvent { // ParkEvent 继承自平台相关的事件实现
    private:
        JavaThread * _AssociatedWith; // 指向与这个 ParkEvent 关联的 Java 线程对象
    	// ...
    public:
        ParkEvent(): _AssociatedWith(NULL) {} // 构造函数

        void associate_with(JavaThread * thread) { // 关联线程
            _AssociatedWith = thread;
        }
        JavaThread * associated_with() {
            return _AssociatedWith;
        }
        // ...
};

// 以 linux 为例
class PlatformEvent: public CHeapObj < mtInternal > {
    private: 
        int _Event; // 0:初识值,1:当前线程被 unpark,-1:当前线程被 park
        int _nParked; // 0:当前线程未被 park,1:当前线程被 park
        pthread_mutex_t _mutex[1]; // 互斥锁
        pthread_cond_t _cond[1]; // 条件变量
        // ...
    public: 
        PlatformEvent();
        ~PlatformEvent();
        void park();
        void park(long millis);
        void unpark();
        // ...
};

​ 我们可以看到 PlatformEvent 果不其然是使用了 Mutex 和 Condition 的结合,那么它究竟是怎样使用的呢?接下来就来看看 park 和 unpark 的源码。

int os::PlatformEvent::park(jlong millis) {
    // 确保当前线程没有被 park
    guarantee(_nParked == 0, "invariant");

    int v;
    for (;;) {
        // 先保存初始 _Event 的值
        v = _Event;
        // 使用 CAS 对 _Event 自减
        if (Atomic::cmpxchg(v - 1, & _Event, v) == v) break;
    }
    guarantee(v >= 0, "invariant");

    // unpark 的时候会设置 _Event = 1,因此如果 v != 0 表示其他线程 unpark 我,我直接返回
    if (v != 0) return OS_OK;

    // 现在 v 一定为 0,那么线程必须阻塞,首先计算绝对时间 abst 用于超时
    struct timespec abst;
    compute_abstime( & abst, millis);
    int ret = OS_TIMEOUT;
    
    // ========== 加 mutex 锁保护 condition ==========
    int status = pthread_mutex_lock(_mutex);
    assert_status(status == 0, status, "mutex_lock");
    
    guarantee(_nParked == 0, "invariant");
    ++_nParked;

    // _Event < 0 说明没有人来 unpark 我
    // 在 while 循环中为了确保,如果不是正常的 unpark 唤醒,线程会继续等待
    while (_Event < 0) {
        // 释放 mutex 并等待 condition,此时线程被阻塞在 condition 的等待队列中,
        // 被唤醒后重新拿到 mutex,才能继续执行
        status = pthread_cond_timedwait(_cond, _mutex, _abstime);

        assert_status(status == 0 || status == EINTR ||
            status == ETIME || status == ETIMEDOUT,
            status, "cond_timedwait");
        
        // 允许伪唤醒
        if (!FilterSpuriousWakeups) break;
        
        // 超时退出
        if (status == ETIME || status == ETIMEDOUT) break; 
    }
    --_nParked;
    
    // 这表示被正常 unpark 唤醒
    if (_Event >= 0) {
        ret = OS_OK;
    }
    
    // 重置为 0
    _Event = 0;
    
    // ========== 释放 mutex 锁 ==========
    status = pthread_mutex_unlock(_mutex);
    assert_status(status == 0, status, "mutex_unlock");
    assert(_nParked == 0, "invariant");
    OrderAccess::fence();
    return ret;
}

​ 看到这相信大家对 park 的实现已经很清楚了,它就是简单使用了 Mutex 保护 Condition 的经典模型,保证了 park 和 unpark 的互斥。

​ 下面再来看看 unpark。

void os::PlatformEvent::unpark() {

    // 使用 CAS 将 _Event 设为 1,返回其旧值
    // 如果其旧值为 -1,说明线程被 park,需要执行后面的逻辑将其唤醒,否则返回
    if (Atomic::xchg(1, & _Event) >= 0) return;

    // ========== 加 mutex 锁保护 condition ==========
    int status = pthread_mutex_lock(_mutex);
    assert_status(status == 0, status, "mutex_lock");
    
    int AnyWaiters = _nParked;
    assert(AnyWaiters == 0 || AnyWaiters == 1, "invariant");
    
    if (AnyWaiters != 0) {
        AnyWaiters = 0;
        // 发送 condition 信号,唤醒目标线程
        // 线程此时被转移至 mutex 等待队列,待下面 mutex 锁被释放,被 park 的线程重新拿到锁
        status = pthread_cond_signal(_cond);
        assert_status(status == 0, status, "cond_signal");
    }
    
    // ========== 释放 mutex 锁 ==========
    status = pthread_mutex_unlock(_mutex);
    assert_status(status == 0, status, "mutex_unlock");
}

​ 现在我觉得大家对 park 和 unpark 的代码流程已经比较清楚了,但是可能对 Mutex 和 Condition 还是有些疑惑,所以我进一步解释一下。

5. Mutex 和 Condition

​ Mutex 和 Condition 是 Linux 提供的最基础的同步基础设施,是一切代码同步操作能够发生的基础。Mutex 用于保护共享数据,Condition 用于等待,这两个维度缺一不可。

​ 下面我们从 Linux 源码中对 Mutex 的实现来进一步分析。

struct mutex {
    atomic_t count; // 1:锁可用, -1:锁不可用,0:锁不可用但没有等待者
    spinlock_t wait_lock; // 自旋锁,为了安全访问 wait_list
    struct list_head wait_list; // 等待队列
};

// 加锁过程
static inline int __sched
__mutex_lock_common(struct mutex * lock, long state, unsigned int subclass, unsigned long ip) {
    // count 的状态变化:
    // 1 -> -1:成功获取锁
    // 0 -> -1 / -1 -> -1:锁已被占用

    // 获取当前进程控制块
    struct task_struct * task = current;

    struct mutex_waiter waiter;
    unsigned int old_val;
    unsigned long flags;

    // 获取保护等待队列的自旋锁
    spin_lock_mutex( & lock -> wait_lock, flags);
    // 将当前任务加入等待队列尾部
    list_add_tail( & waiter.list, & lock -> wait_list);
    waiter.task = task;

    // 用一条汇编指令完成 [count -> -1,返回原值]
    old_val = atomic_xchg( & lock -> count, -1);
    // 如果原值为 1,表示锁可用,直接拿到锁
    if (old_val == 1)
        goto done;

    // 走到这说明锁被占用

    // 进入等待循环
    for (;;) {
        // 再次尝试获取锁
        old_val = atomic_xchg( & lock -> count, -1);
        if (old_val == 1)
            break;

        // 检查是否有信号需要处理
        if (unlikely((state == TASK_INTERRUPTIBLE && signal_pending(task)) ||
                (state == TASK_KILLABLE && fatal_signal_pending(task)))) {
            mutex_remove_waiter(lock, & waiter, task_thread_info(task));
            mutex_release( & lock -> dep_map, 1, ip);
            spin_unlock_mutex( & lock -> wait_lock, flags);
            debug_mutex_free_waiter( & waiter);
            // 返回被信号中断
            return -EINTR;
        }

        // 如果走到这还不能获取锁

        // 设置进程状态为睡眠
        __set_task_state(task, state);
        // 释放自旋锁
        spin_unlock_mutex( & lock -> wait_lock, flags);
        // 主动调度,让出 CPU
        schedule();

        // 被唤醒后重新获取自旋锁
        spin_lock_mutex( & lock -> wait_lock, flags);
    }

    // 已经获取了锁
    done:
        lock_acquired( & lock -> dep_map);
    // 将该任务从等待队列中删除
    mutex_remove_waiter(lock, & waiter, task_thread_info(task));
    debug_mutex_set_owner(lock, task_thread_info(task));
    // 如果等待队列为空将 count 置为 0
    if (likely(list_empty( & lock -> wait_list)))
        atomic_set( & lock -> count, 0);
    spin_unlock_mutex( & lock -> wait_lock, flags);
    debug_mutex_free_waiter( & waiter);
    return 0;
}

static inline void
__mutex_unlock_common_slowpath(atomic_t * lock_count, int nested) {
    // 通过 count 成员地址找到整个 mutex 结构
    struct mutex * lock = container_of(lock_count, struct mutex, count);
    unsigned long flags;

    // 获取等待队列保护锁
    spin_lock_mutex( & lock -> wait_lock, flags);
    mutex_release( & lock -> dep_map, nested, _RET_IP_);
    debug_mutex_unlock(lock);

    // 将 count 设为 1,表示锁可用
    if (__mutex_slowpath_needs_to_unlock())
        atomic_set( & lock -> count, 1);

    // 如果等待队列不为空,获取等待队列中的第一个任务,唤醒该任务
    if (!list_empty( & lock -> wait_list)) {
        struct mutex_waiter * waiter = list_entry(lock -> wait_list.next, struct mutex_waiter, list);
        debug_mutex_wake_waiter(lock, waiter);
        wake_up_process(waiter -> task);
    }

    debug_mutex_clear_owner(lock);

    // 释放自旋锁
    spin_unlock_mutex( & lock -> wait_lock, flags);
}struct mutex {
   atomic_t  count; // 1:锁可用, -1:锁不可用,0:锁不可用但没有等待者
   spinlock_t  wait_lock; // 自旋锁,为了安全访问 wait_list
   struct list_head wait_list; // 等待队列
};

// 加锁过程
static inline int __sched
__mutex_lock_common(struct mutex * lock, long state, unsigned int subclass, unsigned long ip) {
    // count 的状态变化:
    // 1 -> -1:成功获取锁
    // 0 -> -1 / -1 -> -1:锁已被占用
    
    // 获取当前进程控制块
    struct task_struct * task = current;
    
    struct mutex_waiter waiter;
    unsigned int old_val;
    unsigned long flags;

    // 获取保护等待队列的自旋锁
    spin_lock_mutex( & lock -> wait_lock, flags);
    // 将当前任务加入等待队列尾部
    list_add_tail( & waiter.list, & lock -> wait_list);
    waiter.task = task;

    // 用一条汇编指令完成 [count -> -1,返回原值]
    old_val = atomic_xchg( & lock -> count, -1);
    // 如果原值为 1,表示锁可用,直接拿到锁
    if (old_val == 1)
        goto done;
    
    // 走到这说明锁被占用

    // 进入等待循环
    for (;;) {
        // 再次尝试获取锁
        old_val = atomic_xchg( & lock -> count, -1);
        if (old_val == 1)
            break;

        // 检查是否有信号需要处理
        if (unlikely((state == TASK_INTERRUPTIBLE && signal_pending(task)) 
                     || (state == TASK_KILLABLE && fatal_signal_pending(task)))) {
            mutex_remove_waiter(lock, & waiter, task_thread_info(task));
            mutex_release( & lock -> dep_map, 1, ip);
            spin_unlock_mutex( & lock -> wait_lock, flags);
            debug_mutex_free_waiter( & waiter);
            // 返回被信号中断
            return -EINTR;
        }
        
        // 如果走到这还不能获取锁
        
        // 设置进程状态为睡眠
        __set_task_state(task, state);
        // 释放自旋锁
        spin_unlock_mutex( & lock -> wait_lock, flags);
        // 主动调度,让出 CPU
        schedule();
        
        // 被唤醒后重新获取自旋锁
        spin_lock_mutex( & lock -> wait_lock, flags);
    }
    
    // 已经获取了锁
    done:
        lock_acquired( & lock -> dep_map);
    // 将该任务从等待队列中删除
    mutex_remove_waiter(lock, & waiter, task_thread_info(task));
    debug_mutex_set_owner(lock, task_thread_info(task));
    // 如果等待队列为空将 count 置为 0
    if (likely(list_empty( & lock -> wait_list)))
        atomic_set( & lock -> count, 0);
    spin_unlock_mutex( & lock -> wait_lock, flags);
    debug_mutex_free_waiter( & waiter);
    return 0;
}

static inline void
__mutex_unlock_common_slowpath(atomic_t *lock_count, int nested)
{
   // 通过 count 成员地址找到整个 mutex 结构
   struct mutex *lock = container_of(lock_count, struct mutex, count);
   unsigned long flags;
    
   // 获取等待队列保护锁
   spin_lock_mutex(&lock->wait_lock, flags);
   mutex_release(&lock->dep_map, nested, _RET_IP_);
   debug_mutex_unlock(lock);
   
   // 将 count 设为 1,表示锁可用
   if (__mutex_slowpath_needs_to_unlock())
       atomic_set(&lock->count, 1);

   // 如果等待队列不为空,获取等待队列中的第一个任务,唤醒该任务
   if (!list_empty(&lock->wait_list)) {
       struct mutex_waiter *waiter =list_entry(lock->wait_list.next,struct mutex_waiter, list);
       debug_mutex_wake_waiter(lock, waiter);
       wake_up_process(waiter->task);
   }

   debug_mutex_clear_owner(lock);
    
   // 释放自旋锁
   spin_unlock_mutex(&lock->wait_lock, flags);
}

​ 看到这大家应该已经清楚 Mutex 的锁定逻辑了,首先先有一个 wait_lock 为了安全访问 wait_list,然后核心抢锁逻辑其实是依靠一条汇编指令 old_val = atomic_xchg( & lock -> count, -1);,这条汇编指令保证在单个 CPU 时间片内完成,不会被中断,它主要打包了这几步操作:拿到旧值、设置新值、返回旧值。count 状态机告诉我们只有它的值从 1 变为 -1 才算抢锁成功。

​ 如果抢锁没有成功,就执行 schedule(); 主动让出 CPU,这是线程睡眠发生的位置。

​ 另外还可以看出 wait_list 是队列,并且每次都从头部取出元素,这表示 Mutex 是公平锁。

​ 现在来思考一个问题,只使用 Mutex 能实现 park 和 unpark 吗?肯定是不行的,因为 Mutex 只维护 “锁可用 / 不可用” 状态,而我们的 park 和 unpark 是要维护更上层的 _Event 状态,这就必须使用 Condition。

pthread_cond_wait(&qready, &qlock); 这是 Condition 提供的核心操作接口,它的作用是将 [将线程放入条件等待队列,释放 Mutex] 打包成一个原子操作。所以 Condition 是用于等待特定条件的,而 Mutex 只是用于等待锁的,它们两个缺一不可。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐