查缓存没命中就全砸到库上,连接池满了怎么办?我用 CompletableFuture 把重复请求合并了一下
大多数缓存穿透问题,根本不需要引入布隆过滤器或分布式锁——你需要的是一层本地请求合并。
当一个热点 key 缓存失效,几十个并发请求同时打到数据库,连接池瞬间吃满,其他正常请求全部排队超时。这个场景很多团队都遇到过。连接池资源是有限的,比如你配了 maximum-pool-size=20,一旦这 20 个连接全被同一条无效 SQL 占住,整个服务就瘫了。
解决方案其实很直接:同一个 key 的数据库查询,在应用层只允许一个线程执行,其他线程挂起等待结果。Java 里用 CompletableFuture 能很干净地实现这一点。
我先把核心代码写出来,再解释每一层是干什么的。
public class MergingQueryService {
private final ConcurrentHashMap<String, CompletableFuture<User>> pendingQueries
= new ConcurrentHashMap<>();
private final UserRepository userRepository;
public MergingQueryService(UserRepository userRepository) {
this.userRepository = userRepository;
}
public User getUserById(String userId) throws Exception {
// 第1步:尝试注册一个查询
CompletableFuture<User> future = new CompletableFuture<>();
CompletableFuture<User> existingFuture = pendingQueries.putIfAbsent(userId, future);
if (existingFuture != null) {
// 第2步:已有线程在查询,直接等待结果
return existingFuture.get(3, TimeUnit.SECONDS);
}
try {
// 第3步:只有第一个线程会执行到这里
User user = userRepository.findById(userId);
future.complete(user);
return user;
} catch (Exception e) {
future.completeExceptionally(e);
throw e;
} finally {
// 第4步:无论成功失败,都要移除 future
pendingQueries.remove(userId);
}
}
}
这段代码的逻辑很直接。
第1步:每个请求进来,先用 putIfAbsent 尝试往 Map 里放一个空的 CompletableFuture。这个方法返回 null 说明你是第一个线程,返回非 null 说明已经有人在处理了。
第2步:如果你不是第一个线程,拿到的是别人创建的 CompletableFuture,直接调用 get() 挂起当前线程,等第一个线程把结果填进来。这里设置了 3 秒超时,防止第一个线程挂了导致无限等待。
第3步:如果你是第一个线程,正常执行数据库查询。查到结果后调用 future.complete(user),所有挂起的线程会立即被唤醒并拿到同一个结果。
第4步:这是最容易踩坑的地方。finally 里的 remove 至关重要——如果不移除,这个 key 会永远留在 Map 里,后续所有请求都会等一个永远不会被完成的 Future,直到超时。用 remove(userId) 而不是 remove(userId, future),因为万一当前线程在 remove 之前抛了别的异常导致 future 对象引用丢失,用精确匹配可能删不掉。
一个真实的压测数据:在我们一个用户查询服务上,用 JMeter 模拟 100 个并发线程请求同一个不存在的 userId,数据库连接池大小 10。优化前,10 个连接全部被占满,剩下 90 个请求在等连接,平均响应时间飙到 8 秒;加上这层合并后,实际落到数据库的查询只有 1 次,平均响应时间降到 12 毫秒。
但这只是最基础的版本,生产环境还需要处理几个问题。
失败的请求不应该被缓存
上面的代码有个问题:如果查出来的结果是 null(用户不存在),future.complete(null) 也会让等待线程拿到 null。但 null 本身不是异常,调用方需要判断,容易漏掉。
更好的做法是,把"用户不存在"也当作一种正常结果返回,但用 Optional 包装:
public class MergingQueryService {
private final ConcurrentHashMap<String, CompletableFuture<Optional<User>>> pendingQueries
= new ConcurrentHashMap<>();
public Optional<User> getUserById(String userId) throws Exception {
CompletableFuture<Optional<User>> future = new CompletableFuture<>();
CompletableFuture<Optional<User>> existingFuture = pendingQueries.putIfAbsent(userId, future);
if (existingFuture != null) {
return existingFuture.get(3, TimeUnit.SECONDS);
}
try {
User user = userRepository.findById(userId);
// user 可能是 null,但 Optional 本身不是 null
future.complete(Optional.ofNullable(user));
return Optional.ofNullable(user);
} catch (Exception e) {
future.completeExceptionally(e);
throw e;
} finally {
pendingQueries.remove(userId);
}
}
}
注意:不要把异常的 Future 留在 Map 里。completeExceptionally 之后,等待线程调用 get() 会抛出 ExecutionException,但 Future 对象本身还在 Map 里。这就是为什么 finally 里必须 remove——无论成功还是失败,都要清理。
防止 Map 无限增长
ConcurrentHashMap 本身没有淘汰机制。如果 key 的基数很大,比如 userId 有上千万个,即使每个 key 只存在几毫秒,Map 的容量也可能膨胀到影响 GC。
加个上限很简单:
public class BoundedMergingQueryService {
private final int maxPendingQueries;
private final ConcurrentHashMap<String, CompletableFuture<Optional<User>>> pendingQueries;
public BoundedMergingQueryService(int maxPendingQueries, UserRepository userRepository) {
this.maxPendingQueries = maxPendingQueries;
this.pendingQueries = new ConcurrentHashMap<>();
this.userRepository = userRepository;
}
public Optional<User> getUserById(String userId) throws Exception {
// 超过上限时直接降级为独立查询,不合并
if (pendingQueries.size() >= maxPendingQueries) {
return Optional.ofNullable(userRepository.findById(userId));
}
CompletableFuture<Optional<User>> future = new CompletableFuture<>();
CompletableFuture<Optional<User>> existingFuture = pendingQueries.putIfAbsent(userId, future);
if (existingFuture != null) {
try {
return existingFuture.get(3, TimeUnit.SECONDS);
} catch (TimeoutException e) {
// 超时也降级
return Optional.ofNullable(userRepository.findById(userId));
}
}
try {
User user = userRepository.findById(userId);
future.complete(Optional.ofNullable(user));
return Optional.ofNullable(user);
} catch (Exception e) {
future.completeExceptionally(e);
throw e;
} finally {
pendingQueries.remove(userId);
}
}
}
设置一个合理的上限,比如 1000。当并发查询的不同 key 数量超过这个值时,新请求直接走数据库,不再合并。这类似于限流里的"快速失败"策略——宁可多打几次数据库,也不能让内存炸掉。
实际生产环境中,pendingQueries.size() 在 1000 以内时,每个 entry 存活时间通常只有几十毫秒(取决于数据库查询耗时),GC 压力几乎可以忽略。
和本地缓存的配合
请求合并解决的是"同一时刻多个请求查同一个 key"的问题。如果这个 key 隔几秒就被查一次,合并帮不上忙——因为每次都是一个新的查询周期,Map 里没有 pending 的 Future。
这时候需要加一层本地缓存,比如 Caffeine:
public class CachedMergingQueryService {
private final Cache<String, Optional<User>> localCache;
private final ConcurrentHashMap<String, CompletableFuture<Optional<User>>> pendingQueries;
private final UserRepository userRepository;
public CachedMergingQueryService(UserRepository userRepository) {
this.userRepository = userRepository;
this.pendingQueries = new ConcurrentHashMap<>();
this.localCache = Caffeine.newBuilder()
.maximumSize(10_000)
.expireAfterWrite(5, TimeUnit.SECONDS)
.build();
}
public Optional<User> getUserById(String userId) throws Exception {
// 第1层:本地缓存
Optional<User> cached = localCache.getIfPresent(userId);
if (cached != null) {
return cached;
}
// 第2层:请求合并
CompletableFuture<Optional<User>> future = new CompletableFuture<>();
CompletableFuture<Optional<User>> existingFuture = pendingQueries.putIfAbsent(userId, future);
if (existingFuture != null) {
return existingFuture.get(3, TimeUnit.SECONDS);
}
try {
User user = userRepository.findById(userId);
Optional<User> result = Optional.ofNullable(user);
future.complete(result);
// 回填缓存
localCache.put(userId, result);
return result;
} catch (Exception e) {
future.completeExceptionally(e);
throw e;
} finally {
pendingQueries.remove(userId);
}
}
}
这里缓存过期时间设了 5 秒,比 Redis 缓存的 60 秒短很多。本地缓存只是用来消化"缓存刚过期的那一瞬间"的并发流量,真正的缓存权威数据源还是 Redis。5 秒足够让那波并发请求合并成一次数据库查询,之后本地缓存也过期了,下次请求重新走 Redis -> DB 的路径。
如果 Redis 里存的就是 null(表示用户不存在),本地缓存也应该存 Optional.empty(),Caffeine 会正常缓存它,避免反复穿透。
常见问题
能不能直接用 synchronized 或者 Lock?
可以,但效果差很多。synchronized 会让所有请求串行化——不同 userId 的查询本可以并行执行,用一把全局锁就全堵住了。你可以用 ConcurrentHashMap 存 lock 对象,每个 key 一把锁,但手动管理锁的创建和销毁很麻烦,容易内存泄漏。CompletableFuture 的好处是天然异步,第一个线程执行查询时,后续线程只是挂起在 get() 上,不占用额外线程资源。
请求合并后,第一个线程查库耗时很长,等待线程会不会超时?
会,所以要设置合理的超时时间。我习惯设 3 秒,超过就抛 TimeoutException。调用方可以选择重试或降级。注意超时后 get() 会抛异常,但第一个线程还在执行——它最终会 complete 那个 Future,只是没人等了。这就是为什么 finally 里必须 remove,否则这个"孤儿 Future"会一直留在 Map 里。
这个方案和 Hystrix/Sentinel 的请求合并有什么区别?
Hystrix 的请求合并是把一个时间窗口内的多个请求合并成一次批量查询,比如把 10 个单独的 getUserById 合并成一条 SELECT * FROM users WHERE id IN (...)。它解决的是"减少数据库交互次数",但窗口期内请求还是要等。而本文的方案是"同一 key 的重复查询只执行一次",两个维度不同,可以组合使用。