You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现无竞态条件的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后完成,也不会把失效的结果写入缓存。

具体实现:

  1. 添加原子版本号:
    private final AtomicInteger version = new AtomicInteger(0);
    
  2. 修改clear方法:
    public void clear() {
        int newVersion = version.incrementAndGet();
        results.clear(); // 先清结果缓存
        executions.clear(); // 再清执行跟踪,新请求会创建新任务
    }
    
  3. 任务提交时记录版本号,完成时校验:
    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);
    });
    

三、关于键级读写锁方案的优化

你提出的键级读写锁思路是可行的,但有几个细节需要修正,避免锁管理错误:

问题点修正

  1. 不要随意移除锁对象:你在finally块中移除locks里的锁,会导致其他持有该锁的线程出现异常。锁对象是轻量级的,留在ConcurrentHashMap中即可,无需主动清理。
  2. 正确处理锁升级: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 23:52:40