如何实现无竞态条件的Java CompletableFuture原子缓存?
解决AsyncCache的竞态条件问题
咱们一步步拆解你的问题,先解决你发现的两个核心同步问题,再聊聊你提出的键级读写锁方案,最后搞定clear方法的同步逻辑。
一、修复第一个竞态条件:任务完成时的间隙问题
你提到的第一个问题——executions.remove(k)和results.put(k, r)之间的间隙会导致重复执行,核心原因是先移除执行跟踪,再存入结果,这会让新请求误以为没有正在执行的任务,从而触发重复调用。
修复思路
调整操作顺序:执行成功时先把结果存入results,再从executions移除跟踪;执行失败时直接移除跟踪。这样新请求进来时,要么能从results拿到结果,要么能从executions找到正在执行的任务,不会触发重复执行。同时,失败的任务不存入结果,允许后续请求重试,这也符合缓存的预期行为。
修改后的whenComplete逻辑:
return future.whenComplete((r, t) -> { if (t == null) { // 仅成功时存入结果 results.put(k, r); } executions.remove(k); // 无论成败,都移除执行跟踪 });
二、修复第二个竞态条件:clear方法的间隙问题
你发现的第二个问题——results.clear()和executions.clear()之间的间隙,会导致新请求拿到已失效的任务结果。核心原因是两次clear不是原子操作,旧任务的结果可能在clear完成后被错误存入缓存。
修复思路
引入全局版本号,每次调用clear时递增版本号;任务提交时记录当前版本号,任务完成后只有版本号与当前一致,才允许将结果存入results。这样即使旧任务在clear后完成,也不会把失效的结果写入缓存。
具体实现:
- 添加原子版本号:
private final AtomicInteger version = new AtomicInteger(0); - 修改
clear方法:public void clear() { int newVersion = version.incrementAndGet(); results.clear(); // 先清结果缓存 executions.clear(); // 再清执行跟踪,新请求会创建新任务 } - 任务提交时记录版本号,完成时校验:
final int taskVersion = version.get(); CompletableFuture<Object> future = CompletableFuture.supplyAsync(() -> { try { return callable.call(); } catch (Exception e) { throw new CompletionException(e); } }, executor); return future.whenComplete((r, t) -> { if (t == null && version.get() == taskVersion) { // 版本一致且成功才存入 results.put(k, r); } executions.remove(k); });
三、关于键级读写锁方案的优化
你提出的键级读写锁思路是可行的,但有几个细节需要修正,避免锁管理错误:
问题点修正
- 不要随意移除锁对象:你在
finally块中移除locks里的锁,会导致其他持有该锁的线程出现异常。锁对象是轻量级的,留在ConcurrentHashMap中即可,无需主动清理。 - 正确处理锁升级:
ReentrantReadWriteLock不支持直接从读锁升级为写锁,需要先释放读锁,再获取写锁,同时要做双重检查(防止其他线程在间隙中写入结果)。
优化后的键级锁版本代码
import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.locks.ReentrantReadWriteLock; public class AsyncCache { private final ExecutorService executor = Executors.newFixedThreadPool(10); private final ConcurrentHashMap<String, CompletableFuture<Object>> executions = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, Object> results = new ConcurrentHashMap<>(); private final ConcurrentMap<String, ReentrantReadWriteLock> keyLocks = new ConcurrentHashMap<>(); private final AtomicInteger version = new AtomicInteger(0); public CompletableFuture<Object> get(String key, Callable<Object> callable) { ReentrantReadWriteLock lock = keyLocks.computeIfAbsent(key, k -> new ReentrantReadWriteLock()); boolean readLocked = true; lock.readLock().lock(); try { Object result = results.get(key); if (result != null) { return CompletableFuture.completedFuture(result); } // 升级为写锁:先释放读锁 lock.readLock().unlock(); readLocked = false; lock.writeLock().lock(); try { // 双重检查,防止其他线程在间隙中写入结果 result = results.get(key); if (result != null) { return CompletableFuture.completedFuture(result); } final int taskVersion = version.get(); return executions.computeIfAbsent(key, k -> { CompletableFuture<Object> future = CompletableFuture.supplyAsync(() -> { try { return callable.call(); } catch (Exception e) { throw new CompletionException(e); } }, executor); return future.whenComplete((r, t) -> { lock.writeLock().lock(); try { if (t == null && version.get() == taskVersion) { results.put(k, r); } executions.remove(k); } finally { lock.writeLock().unlock(); } }); }); } finally { lock.writeLock().unlock(); } } finally { if (readLocked) { lock.readLock().unlock(); } } } public void clear() { int newVersion = version.incrementAndGet(); results.clear(); executions.clear(); } // 关闭线程池,避免资源泄漏 public void shutdown() { executor.shutdown(); } }
这个版本既保留了键级锁的细粒度控制,又通过版本号解决了clear方法的同步问题,同时避免了锁管理错误。
内容的提问来源于stack exchange,提问作者Ricardo Marimon
相关产品推荐
相关产品推荐

