在实际并发编程中,读写锁(Reader-Writer Lock)是解决“读多写少”场景下性能瓶颈的经典工具。标准的读写锁实现,如 Java 中的ReentrantReadWriteLock,通常采用“读优先”或“写优先”的策略,但这两种策略都存在一个潜在问题:在持续高并发请求下,某些线程可能会因为策略原因长时间无法获取锁,即“饥饿”(Starvation)。例如,在写优先策略下,如果写请求源源不断,读线程可能永远无法执行;反之,在读优先策略下,写线程也可能被无限期推迟。
FairRWLock正是为了解决这种公平性问题而设计的一种“防饥饿”的读写锁。它的核心目标是在保证读写操作正确性的前提下,引入一种公平的调度机制,确保无论是读线程还是写线程,在长时间运行中都能获得执行机会,避免任何一方被饿死。这对于需要保证服务响应质量、避免长尾延迟的系统(如实时数据处理、在线交易系统)至关重要。本文将深入探讨公平读写锁的设计原理,并提供一个从零实现的、可运行的 Java 示例,涵盖其核心机制、关键代码、使用方式以及生产环境下的考量。
1. 理解读写锁的公平性问题与 FairRWLock 设计目标
在深入实现之前,必须清晰理解标准读写锁的“不公平”是如何产生的,以及FairRWLock要达成的目标。
1.1 标准读写锁的策略与饥饿风险
常见的读写锁实现基于一个状态变量(通常是一个int)和两个等待队列。状态变量记录当前持有读锁的数量和写锁的持有状态。其调度策略大致分为两类:
- 读优先:只要当前没有写锁被持有,新来的读请求可以立即获取锁,即使有写线程正在等待。这可能导致写线程饥饿。
- 写优先:一旦有写线程开始等待,后续新来的读请求会被阻塞,直到所有等待的写线程完成。这可能导致读线程饥饿。
这两种策略的“优先”都是非公平的,它牺牲了部分线程的等待时间以换取整体吞吐量。但在对延迟敏感或要求严格公平性的场景下,这种牺牲是不可接受的。
1.2 FairRWLock 的核心设计思想
FairRWLock的设计思想是“先到先服务”(FIFO)的公平性。它维护一个统一的等待队列,所有请求锁的线程(无论是读是写)都按照到达顺序排队。但单纯的 FIFO 会严重损害读写锁的并发性(即允许多个读线程同时执行的优势)。因此,FairRWLock需要在公平性和并发性之间做出精巧的平衡。
一个典型的公平读写锁实现遵循以下规则:
- 队列头部线程规则:锁的授予决策主要基于等待队列头部的线程。
- 读线程共享规则:如果队列头部是一个读线程,那么它不仅自己可以获得锁,它之后连续排队的读线程也可以“搭便车”一起获得锁,直到遇到一个写线程为止。这保留了读并发性。
- 写线程独占规则:如果队列头部是一个写线程,那么它独占锁,它之后的线程必须等待。
这种设计确保了:1) 线程按照请求顺序被服务(公平);2) 连续的读操作可以并发执行(高效);3) 写操作不会被后来的读操作无限插队(防写饥饿)。
2. 环境准备与项目结构
我们将使用 Java 语言实现一个简易但功能完整的FairRWLock。选择 Java 是因为其内置的AbstractQueuedSynchronizer (AQS)框架为构建自定义同步器提供了强大且安全的基础。
2.1 环境要求
- JDK 版本:JDK 8 或更高版本(主要使用
AQS和ReentrantLock)。 - 构建工具:Maven 或 Gradle 均可,本文使用 Maven 示例,但核心代码不依赖特定构建工具。
- IDE:任何支持 Java 的 IDE(如 IntelliJ IDEA, Eclipse)或文本编辑器。
2.2 项目结构与依赖
创建一个标准的 Maven 项目。由于我们完全基于 JDK 内置并发库,无需额外第三方依赖。
fair-rwlock-demo/ ├── pom.xml └── src/ └── main/ └── java/ └── com/ └── example/ └── concurrent/ ├── FairRWLock.java // 公平读写锁核心实现 └── FairRWLockExample.java // 使用示例和测试pom.xml文件只需最基本的配置:
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>fair-rwlock-demo</artifactId> <version>1.0-SNAPSHOT</version> <properties> <maven.compiler.source>8</maven.compiler.source> <maven.compiler.target>8</maven.compiler.target> </properties> </project>3. 基于 AQS 实现 FairRWLock 核心逻辑
AbstractQueuedSynchronizer (AQS)是构建锁和同步器的框架。它内部维护了一个 FIFO 队列(正是我们需要的公平队列)和状态管理。我们的FairRWLock将作为AQS的一个内部类来实现。
3.1 锁状态定义与读写计数
读写锁的状态需要同时表示读锁数量(可大于1)和写锁状态(0或1)。一个常见的技巧是使用一个 32 位的int,高 16 位表示读锁计数,低 16 位表示写锁计数。
package com.example.concurrent; import java.util.concurrent.locks.AbstractQueuedSynchronizer; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; public class FairRWLock { // 内部同步器,继承AQS private final Sync sync = new Sync(); // 读锁和写锁的视图 private final Lock readLock = new ReadLock(); private final Lock writeLock = new WriteLock(); public Lock readLock() { return readLock; } public Lock writeLock() { return writeLock; } /** * 同步控制核心。使用一个int state表示锁状态。 * 高16位:读锁持有次数 (readCount) * 低16位:写锁持有次数 (writeCount, 0或1) 和写锁重入次数 * 由于是公平锁,尝试获取的逻辑完全由AQS的队列机制决定。 */ private static final class Sync extends AbstractQueuedSynchronizer { private static final long serialVersionUID = 1L; // 位偏移常量 static final int SHARED_SHIFT = 16; static final int SHARED_UNIT = (1 << SHARED_SHIFT); static final int MAX_COUNT = (1 << SHARED_SHIFT) - 1; // 65535 static final int EXCLUSIVE_MASK = (1 << SHARED_SHIFT) - 1; // 65535 // 读锁计数 = state >>> 16 (无符号右移) static int sharedCount(int c) { return c >>> SHARED_SHIFT; } // 写锁计数 = state & 65535 static int exclusiveCount(int c) { return c & EXCLUSIVE_MASK; } // 每个读线程的重入计数由ThreadLocal维护 private transient ThreadLocalHoldCounter readHolds; // 缓存最后一个成功获取读锁的线程的计数,优化性能 private transient HoldCounter cachedHoldCounter; Sync() { readHolds = new ThreadLocalHoldCounter(); setState(getState()); // 触发readHolds初始化 } // HoldCounter 用于记录每个线程持有读锁的次数 static final class HoldCounter { int count = 0; final long tid = Thread.currentThread().getId(); } static final class ThreadLocalHoldCounter extends ThreadLocal<HoldCounter> { public HoldCounter initialValue() { return new HoldCounter(); } }3.2 尝试获取读锁(tryAcquireShared)
在公平锁中,读锁的获取必须检查队列。只有当自己是队列头节点,或者队列头节点也是读锁请求时,才能尝试获取。
// 尝试以共享模式获取(读锁) @Override protected int tryAcquireShared(int unused) { Thread current = Thread.currentThread(); int c = getState(); // 如果已经有写锁被持有,且持有者不是当前线程,则失败 if (exclusiveCount(c) != 0 && getExclusiveOwnerThread() != current) { return -1; } int r = sharedCount(c); // 当前读锁数量 // 关键:检查是否应该阻塞。在公平模式下,如果队列中有等待者,且自己不是头节点后的连续读线程,则应阻塞。 if (!readerShouldBlock() && r < MAX_COUNT) { // 使用CAS增加读锁计数 if (compareAndSetState(c, c + SHARED_UNIT)) { // CAS成功,更新线程本地计数 if (r == 0) { firstReader = current; firstReaderHoldCount = 1; } else if (firstReader == current) { firstReaderHoldCount++; } else { HoldCounter rh = cachedHoldCounter; if (rh == null || rh.tid != current.getId()) cachedHoldCounter = rh = readHolds.get(); else if (rh.count == 0) readHolds.set(rh); rh.count++; } return 1; // 获取成功 } } // 如果应该阻塞或CAS失败,进入完整版本的获取逻辑(包含重试和排队) return fullTryAcquireShared(current); } // 决定读线程是否应该阻塞。公平锁的实现:检查队列中是否有其他等待者,且自己不是“合法”的后续读线程。 final boolean readerShouldBlock() { // hasQueuedPredecessors() 是AQS方法,检查当前线程前是否有其他线程在排队。 // 这是公平性的核心:只要有前辈在等,我就不能插队。 return hasQueuedPredecessors(); } // 完整版的共享获取,用于处理重试和重入 final int fullTryAcquireShared(Thread current) { HoldCounter rh = null; for (;;) { int c = getState(); if (exclusiveCount(c) != 0) { if (getExclusiveOwnerThread() != current) return -1; } else if (readerShouldBlock()) { // 再次检查是否应该阻塞 if (firstReader == current) { // 第一个读线程重入,允许 } else { if (rh == null) { rh = cachedHoldCounter; if (rh == null || rh.tid != current.getId()) { rh = readHolds.get(); } } if (rh.count == 0) { // 该线程读锁计数为0,说明是新的读请求,且需要阻塞 readHolds.remove(); return -1; } } } if (sharedCount(c) == MAX_COUNT) throw new Error("Maximum lock count exceeded"); if (compareAndSetState(c, c + SHARED_UNIT)) { // 成功获取,更新计数(类似tryAcquireShared中的逻辑) if (sharedCount(c) == 0) { firstReader = current; firstReaderHoldCount = 1; } else if (firstReader == current) { firstReaderHoldCount++; } else { if (rh == null) rh = cachedHoldCounter; if (rh == null || rh.tid != current.getId()) rh = readHolds.get(); else if (rh.count == 0) readHolds.set(rh); rh.count++; cachedHoldCounter = rh; // 缓存 } return 1; } } }3.3 尝试获取写锁(tryAcquire)
写锁是独占的。在公平锁中,只要队列中有其他等待者(hasQueuedPredecessors()返回true),当前线程就不能获取写锁,必须排队。
// 尝试以独占模式获取(写锁) @Override protected boolean tryAcquire(int acquires) { Thread current = Thread.currentThread(); int c = getState(); int w = exclusiveCount(c); if (c != 0) { // 状态不为0,说明有锁被持有 // 情况1: 有读锁 (w == 0) -> 失败 // 情况2: 有写锁,但持有者不是当前线程 -> 失败 if (w == 0 || current != getExclusiveOwnerThread()) return false; // 情况3: 有写锁,且是当前线程重入 -> 检查是否超限 if (w + exclusiveCount(acquires) > MAX_COUNT) throw new Error("Maximum lock count exceeded"); // 重入,直接设置状态 setState(c + acquires); return true; } // 状态为0,无线程持有锁 // 公平性核心:如果有线程排队在自己前面,就放弃获取,去排队。 if (writerShouldBlock() || !compareAndSetState(c, c + acquires)) return false; // CAS成功,设置独占所有者 setExclusiveOwnerThread(current); return true; } // 决定写线程是否应该阻塞。公平锁:检查是否有前辈在排队。 final boolean writerShouldBlock() { return hasQueuedPredecessors(); }3.4 尝试释放锁(tryReleaseShared 与 tryRelease)
释放逻辑相对直接,主要是对状态进行安全的减少操作。
// 尝试释放共享锁(读锁) @Override protected boolean tryReleaseShared(int unused) { Thread current = Thread.currentThread(); // 更新线程本地持有计数 if (firstReader == current) { if (firstReaderHoldCount == 1) firstReader = null; else firstReaderHoldCount--; } else { HoldCounter rh = cachedHoldCounter; if (rh == null || rh.tid != current.getId()) rh = readHolds.get(); int count = rh.count; if (count <= 1) { readHolds.remove(); if (count <= 0) throw unmatchedUnlockException(); } --rh.count; } // 循环CAS减少state中的读计数 for (;;) { int c = getState(); int nextc = c - SHARED_UNIT; if (compareAndSetState(c, nextc)) // 释放读锁对等待线程的唤醒,由AQS框架处理 return nextc == 0; // 返回true表示此次释放后读锁完全空闲 } } // 尝试释放独占锁(写锁) @Override protected boolean tryRelease(int releases) { if (!isHeldExclusively()) throw new IllegalMonitorStateException(); int nextc = getState() - releases; boolean free = exclusiveCount(nextc) == 0; if (free) { setExclusiveOwnerThread(null); } setState(nextc); return free; } @Override protected boolean isHeldExclusively() { return getExclusiveOwnerThread() == Thread.currentThread(); } // 条件变量支持(可选,但通常需要) final ConditionObject newCondition() { return new ConditionObject(); } } // 读锁视图实现 private final class ReadLock implements Lock { public void lock() { sync.acquireShared(1); } public void lockInterruptibly() throws InterruptedException { sync.acquireSharedInterruptibly(1); } public boolean tryLock() { return sync.tryAcquireShared(1) >= 0; } public boolean tryLock(long timeout, TimeUnit unit) throws InterruptedException { return sync.tryAcquireSharedNanos(1, unit.toNanos(timeout)); } public void unlock() { sync.releaseShared(1); } public Condition newCondition() { throw new UnsupportedOperationException(); } } // 写锁视图实现 private final class WriteLock implements Lock { public void lock() { sync.acquire(1); } public void lockInterruptibly() throws InterruptedException { sync.acquireInterruptibly(1); } public boolean tryLock() { return sync.tryAcquire(1); } public boolean tryLock(long timeout, TimeUnit unit) throws InterruptedException { return sync.tryAcquireNanos(1, unit.toNanos(timeout)); } public void unlock() { sync.release(1); } public Condition newCondition() { return sync.newCondition(); } } }4. 运行验证与行为演示
为了验证FairRWLock的公平性,我们编写一个测试程序,模拟读写混合的场景,并观察线程的执行顺序。
4.1 创建测试用例
package com.example.concurrent; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; public class FairRWLockExample { private static final FairRWLock lock = new FairRWLock(); private static String sharedResource = "Initial"; private static final AtomicInteger readCount = new AtomicInteger(0); private static final AtomicInteger writeCount = new AtomicInteger(0); static class Reader implements Runnable { private final String name; private final CountDownLatch startLatch; Reader(String name, CountDownLatch startLatch) { this.name = name; this.startLatch = startLatch; } @Override public void run() { try { startLatch.await(); // 等待同时开始 lock.readLock().lock(); try { int rc = readCount.incrementAndGet(); System.out.printf("[Time: %3dms] Reader %s STARTED reading. Content: %s (Concurrent Readers: %d)%n", System.currentTimeMillis() % 1000, name, sharedResource, rc); Thread.sleep(100); // 模拟读操作耗时 System.out.printf("[Time: %3dms] Reader %s FINISHED reading.%n", System.currentTimeMillis() % 1000, name); } finally { lock.readLock().unlock(); readCount.decrementAndGet(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } static class Writer implements Runnable { private final String name; private final String newValue; private final CountDownLatch startLatch; Writer(String name, String newValue, CountDownLatch startLatch) { this.name = name; this.newValue = newValue; this.startLatch = startLatch; } @Override public void run() { try { startLatch.await(); lock.writeLock().lock(); try { int wc = writeCount.incrementAndGet(); System.out.printf("[Time: %3dms] Writer %s STARTED writing. New Value: %s%n", System.currentTimeMillis() % 1000, name, newValue); Thread.sleep(200); // 模拟写操作耗时,比读长 sharedResource = newValue; System.out.printf("[Time: %3dms] Writer %s FINISHED writing.%n", System.currentTimeMillis() % 1000, name); } finally { lock.writeLock().unlock(); writeCount.decrementAndGet(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } public static void main(String[] args) throws InterruptedException { System.out.println("=== 测试 FairRWLock 公平性 ==="); System.out.println("启动顺序:R1, R2, W1, R3, W2"); System.out.println("预期行为:R1和R2并发执行,然后W1,然后R3,最后W2。"); System.out.println("------------------------------"); CountDownLatch startLatch = new CountDownLatch(1); Thread r1 = new Thread(new Reader("R1", startLatch)); Thread r2 = new Thread(new Reader("R2", startLatch)); Thread w1 = new Thread(new Writer("W1", "Written_by_W1", startLatch)); Thread r3 = new Thread(new Reader("R3", startLatch)); Thread w2 = new Thread(new Writer("W2", "Written_by_W2", startLatch)); // 按顺序启动线程,模拟请求到达顺序 r1.start(); Thread.sleep(10); // 微小间隔,确保启动顺序 r2.start(); Thread.sleep(10); w1.start(); Thread.sleep(10); r3.start(); Thread.sleep(10); w2.start(); Thread.sleep(50); // 确保所有线程都已启动并进入等待 startLatch.countDown(); // 同时释放所有线程去争抢锁 // 等待所有线程结束 r1.join(); r2.join(); w1.join(); r3.join(); w2.join(); System.out.println("=== 测试结束 ==="); } }4.2 分析运行结果与预期
运行上述程序,观察控制台输出。一个典型的、体现公平性的输出可能如下:
=== 测试 FairRWLock 公平性 === 启动顺序:R1, R2, W1, R3, W2 预期行为:R1和R2并发执行,然后W1,然后R3,最后W2。 ------------------------------ [Time: 123ms] Reader R1 STARTED reading. Content: Initial (Concurrent Readers: 1) [Time: 123ms] Reader R2 STARTED reading. Content: Initial (Concurrent Readers: 2) [Time: 224ms] Reader R1 FINISHED reading. [Time: 224ms] Reader R2 FINISHED reading. [Time: 225ms] Writer W1 STARTED writing. New Value: Written_by_W1 [Time: 426ms] Writer W1 FINISHED writing. [Time: 427ms] Reader R3 STARTED reading. Content: Written_by_W1 (Concurrent Readers: 1) [Time: 528ms] Reader R3 FINISHED reading. [Time: 529ms] Writer W2 STARTED writing. New Value: Written_by_W2 [Time: 730ms] Writer W2 FINISHED writing. === 测试结束 ===结果分析:
- R1 和 R2 并发:它们最先启动,且都是读请求,因此作为队列头部连续的读线程,同时获取锁并执行。
- W1 紧随其后:在 R1/R2 释放锁后,队列头部的线程是 W1(写),因此 W1 获得锁。注意,尽管在 W1 执行期间 R3 早已就绪,但它必须等待。
- R3 在 W1 之后:W1 释放锁后,队列头部的线程是 R3(读),因此 R3 获得锁。
- W2 最后:R3 释放锁后,队列头部的线程是 W2(写),因此 W2 获得锁并执行。
这个顺序严格遵循了 FIFO 公平性,同时允许了读并发(R1和R2)。如果使用非公平锁(如ReentrantReadWriteLock的非公平模式),输出顺序可能是不可预测的,例如 W1 或 W2 可能会被延迟更久。
5. 常见问题排查与性能考量
在实际使用自定义的FairRWLock时,可能会遇到一些典型问题。
5.1 常见问题与排查表
| 问题现象 | 可能原因 | 检查点与解决方案 |
|---|---|---|
| 死锁 | 1. 同一个线程先获取读锁,再尝试获取写锁(锁升级)。 2. 多个线程以不同的顺序获取读写锁。 | 1.FairRWLock不支持锁升级。确保线程在释放读锁前,不要尝试获取写锁。如果需要,请使用支持升级的锁(如 StampedLock)。2. 检查代码中锁的获取顺序,确保全局一致的锁顺序。 |
| 性能低于非公平锁 | 在高争用场景下,严格的 FIFO 排队会导致更多的上下文切换和更低的吞吐量。 | 这是公平锁的设计取舍。如果系统吞吐量优先且可以接受饥饿,应换用非公平锁。使用性能分析工具(如 JProfiler, Async Profiler)确认瓶颈是否确实在锁排队上。 |
IllegalMonitorStateException | 1. 未配对调用lock()/unlock()。2. 一个线程释放了另一个线程持有的锁。 | 1. 确保每个lock()都有对应的unlock(),且放在finally块中。2. 检查锁对象的作用域和线程持有关系。 |
| 读锁重入计数错误 | 线程本地存储ThreadLocalHoldCounter未正确清理,导致内存泄漏或计数错误。 | 确保tryReleaseShared中当线程读计数降为0时,调用readHolds.remove()。我们的实现已包含此逻辑。 |
| 长时间等待 | 队列前有一个持有锁时间很长的写线程,或有一批连续的读线程。 | 这是公平锁的正常行为。检查业务逻辑: 1. 写操作是否可以优化(减小临界区)? 2. 读操作是否必要持有锁?考虑使用乐观锁或拷贝数据。 |
5.2 性能与适用场景权衡
FairRWLock通过引入严格的排队机制解决了饥饿问题,但也付出了性能代价:
- 优点:
- 严格的公平性:完全杜绝饥饿,提供可预测的等待时间。
- 避免长尾延迟:对延迟敏感型应用友好。
- 缺点:
- 吞吐量降低:线程切换和排队开销增加,在高争用下吞吐量通常低于非公平锁。
- 实现复杂度高:需要精细控制队列和状态。
选型建议:
- 使用
FairRWLock:当系统需要保证所有请求(无论读写)都能在有限时间内得到处理时,如实时竞价系统、硬件控制、某些游戏服务器。 - 使用标准非公平
ReentrantReadWriteLock:当系统吞吐量是首要目标,且偶尔的线程饥饿可以接受时,如大多数 Web 应用后端缓存。 - 考虑其他选择:
StampedLock:提供了乐观读、锁升级等更丰富的操作,在特定读远多于写的场景性能可能更好,但 API 更复杂。- 无锁数据结构:彻底避免锁争用,但实现难度极高。
6. 生产环境最佳实践与扩展方向
将FairRWLock用于生产环境,需要考虑更多维度的健壮性。
6.1 监控与诊断
- 队列长度监控:可以通过扩展
Sync类,暴露getQueueLength()方法(AQS 已提供)来监控等待队列长度。队列持续过长是系统过载或锁竞争激烈的信号。 - 持有时间监控:在
lock()和unlock()前后记录时间戳,统计锁持有时间。过长的持有时间会直接影响系统吞吐量和延迟。 - 线程转储分析:当发生疑似死锁或性能问题时,使用
jstack或ThreadMXBean获取线程转储,检查哪些线程阻塞在FairRWLock上。
6.2 代码健壮性增强
- 添加单元测试:覆盖并发场景,如读写交错、重入、中断、超时等。使用
JUnit和ConcurrentUnit等库。 - 实现
toString():重写FairRWLock的toString()方法,输出当前读锁数、写锁持有者、等待队列长度等信息,便于调试。 - 考虑序列化:如果锁对象需要序列化(通常不推荐),需要将
transient字段(如readHolds)妥善处理。
6.3 扩展方向
- 支持锁降级:允许持有写锁的线程同时获取读锁,然后释放写锁,从而降级为读锁。这需要在
tryAcquireShared中检查当前线程是否为写锁持有者。 - 实现可中断的公平策略:在某些场景下,可能希望等待时间过长的读请求可以“让位”给写请求,反之亦然。这需要更复杂的队列管理和策略判断。
- 与
CompletableFuture集成:包装lock()和unlock()操作,使其能够更好地融入异步编程范式。
公平读写锁是并发工具箱中一把精密而特定的工具。理解其实现机制,能帮助开发者在公平与性能之间做出明智的架构选择,并能在遇到相关问题时进行有效排查。本文提供的实现是一个教学示例,揭示了其核心原理。在生产中使用时,建议优先考虑经过充分测试的库(如 JDK 自身的ReentrantReadWriteLock公平模式),或在充分评估后基于此模式进行定制化增强。