尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Java公平读写锁FairRWLock实现:防饥饿与FIFO调度原理详解

Java公平读写锁FairRWLock实现:防饥饿与FIFO调度原理详解 在实际并发编程中读写锁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 version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion groupIdcom.example/groupId artifactIdfair-rwlock-demo/artifactId version1.0-SNAPSHOT/version properties maven.compiler.source8/maven.compiler.source maven.compiler.target8/maven.compiler.target /properties /project3. 基于 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 ThreadLocalHoldCounter { 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确认瓶颈是否确实在锁排队上。IllegalMonitorStateException1. 未配对调用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公平模式或在充分评估后基于此模式进行定制化增强。
返回列表