连接池核心 - ConcurrentBag
HikariCP ConcurrentBag 源码分析
连接池的本质问题:多线程高频地”借一个连接、用完还回来”,传统实现(如 commons-dbcp 用一把大锁守护空闲队列)在高并发下锁竞争成为瓶颈。
HikariCP 的答案是 ConcurrentBag——一个专为连接池定制的并发容器:借出走三级查找(ThreadLocal 缓存 → 共享列表 CAS 扫描 → handoffQueue 等待),常见路径完全无锁,状态流转只靠一次 CAS。
// 基于本地 HikariCP 仓 (JavaSourceReadingLab/frameworks/hikaricp), com/zaxxer/hikari/HikariDataSource.java
public Connection getConnection() throws SQLException
{
if (isClosed()) {
throw new SQLException("HikariDataSource " + this + " has been closed.");
}
if (fastPathPool != null) { // 用 HikariConfig 构造时池在构造器里就建好,走这条免判空快路径
return fastPathPool.getConnection();
}
// See http://en.wikipedia.org/wiki/Double-checked_locking#Usage_in_Java
HikariPool result = pool; // pool 是 volatile,标准双检锁:无参构造+首次 getConnection 才懒初始化
if (result == null) {
synchronized (this) {
result = pool;
if (result == null) {
validate();
LOGGER.info("{} - Starting...", getPoolName());
try {
pool = result = new HikariPool(this);
this.seal(); // 池启动后配置封印,运行期改配置直接抛异常
}
// ... catch PoolInitializationException 略
LOGGER.info("{} - Start completed.", getPoolName());
}
}
}
return result.getConnection();
}// 基于本地 HikariCP 仓 (JavaSourceReadingLab/frameworks/hikaricp), com/zaxxer/hikari/pool/HikariPool.java
public Connection getConnection() throws SQLException
{
return getConnection(connectionTimeout); // 默认 30 秒的 connectionTimeout 就是从这传进去的
}
public Connection getConnection(final long hardTimeout) throws SQLException
{
suspendResumeLock.acquire(); // 池挂起功能的开关;默认不启用时是 FAUX_LOCK 空实现,零开销
final long startTime = currentTime();
try {
long timeout = hardTimeout;
do { // 循环:borrow 到的连接可能是坏的,坏了就关掉再借,直到超时
PoolEntry poolEntry = connectionBag.borrow(timeout, MILLISECONDS);
if (poolEntry == null) {
break; // We timed out... break and throw exception
}
final long now = currentTime();
if (poolEntry.isMarkedEvicted() || (elapsedMillis(poolEntry.lastAccessed, now) > aliveBypassWindowMs && !isConnectionAlive(poolEntry.connection))) {
// 500ms 内刚用过的连接跳过存活检测(aliveBypassWindowMs),避免每次借出都发一条验证 SQL
closeConnection(poolEntry, poolEntry.isMarkedEvicted() ? EVICTED_CONNECTION_MESSAGE : DEAD_CONNECTION_MESSAGE);
timeout = hardTimeout - elapsedMillis(startTime); // 剩余时间接着借,坏连接不吃掉整个超时预算
}
else {
metricsTracker.recordBorrowStats(poolEntry, startTime);
// 泄漏检测就这一句:借出时调度一个延时任务,超过 leakDetectionThreshold 未归还就打警告日志
return poolEntry.createProxyConnection(leakTaskFactory.schedule(poolEntry), now);
}
} while (timeout > 0L);
metricsTracker.recordBorrowTimeoutStats(startTime);
throw createTimeoutException(startTime);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new SQLException(poolName + " - Interrupted during connection acquisition", e);
}
finally {
suspendResumeLock.release();
}
}
// 归还入口:ProxyConnection.close() 最终走到这里,连接不是真关,回到 bag
void recycle(final PoolEntry poolEntry)
{
metricsTracker.recordConnectionUsage(poolEntry);
connectionBag.requite(poolEntry);
}// 基于本地 HikariCP 仓 (JavaSourceReadingLab/frameworks/hikaricp), com/zaxxer/hikari/util/ConcurrentBag.java
public class ConcurrentBag<T extends IConcurrentBagEntry> implements AutoCloseable
{
private final CopyOnWriteArrayList<T> sharedList; // 全量连接都在这,借出也不移除;读无锁,只有增删连接才复制
private final boolean weakThreadLocals;
private final ThreadLocal<List<Object>> threadList; // 每线程的"私藏"缓存,借还的最快路径
private final IBagStateListener listener; // 就是 HikariPool,借不到时通过它异步补建连接
private final AtomicInteger waiters; // 当前有多少线程在等连接
private volatile boolean closed;
private final SynchronousQueue<T> handoffQueue; // 零容量交接队列:归还线程直接把连接"手递手"给等待线程
public interface IConcurrentBagEntry
{
int STATE_NOT_IN_USE = 0;
int STATE_IN_USE = 1;
int STATE_REMOVED = -1;
int STATE_RESERVED = -2; // 房管线程(如 maxLifetime 到期驱逐)先 reserve 占住,防止操作时被借走
boolean compareAndSet(int expectState, int newState); // PoolEntry 用 AtomicIntegerFieldUpdater 实现
void setState(int newState);
int getState();
}
public ConcurrentBag(final IBagStateListener listener)
{
this.listener = listener;
this.weakThreadLocals = useWeakThreadLocals(); // 有自定义 ClassLoader(如 Web 容器热部署)才用弱引用,防内存泄漏
this.handoffQueue = new SynchronousQueue<>(true); // fair=true:等最久的线程先拿到,避免请求饿死
this.waiters = new AtomicInteger();
this.sharedList = new CopyOnWriteArrayList<>();
if (weakThreadLocals) {
this.threadList = ThreadLocal.withInitial(() -> new ArrayList<>(16));
}
else {
this.threadList = ThreadLocal.withInitial(() -> new FastList<>(IConcurrentBagEntry.class, 16));
}
}
public T borrow(long timeout, final TimeUnit timeUnit) throws InterruptedException
{
// Try the thread-local list first
final List<Object> list = threadList.get();
for (int i = list.size() - 1; i >= 0; i--) { // 倒序拿:从尾部删不用挪数组,且最近还的连接最"热"
final Object entry = list.remove(i);
@SuppressWarnings("unchecked")
final T bagEntry = weakThreadLocals ? ((WeakReference<T>) entry).get() : (T) entry;
if (bagEntry != null && bagEntry.compareAndSet(STATE_NOT_IN_USE, STATE_IN_USE)) {
return bagEntry; // 一级命中:全程无锁无竞争。CAS 仍必须做——缓存的连接可能已被别的线程从 sharedList 借走
}
}
// Otherwise, scan the shared list ... then poll the handoff queue
final int waiting = waiters.incrementAndGet();
try {
for (T bagEntry : sharedList) {
if (bagEntry.compareAndSet(STATE_NOT_IN_USE, STATE_IN_USE)) {
// If we may have stolen another waiter's connection, request another bag add.
if (waiting > 1) {
listener.addBagItem(waiting - 1); // 我可能抢了别的等待者的连接,补偿性地让池再建一个
}
return bagEntry; // 二级命中:遍历+CAS,"偷"任何空闲连接,包括别人 ThreadLocal 里缓存的
}
}
listener.addBagItem(waiting); // 真没有了,异步触发建新连接(受 maximumPoolSize 约束)
timeout = timeUnit.toNanos(timeout);
do {
final long start = currentTime();
final T bagEntry = handoffQueue.poll(timeout, NANOSECONDS); // 三级:阻塞等别人归还或新建
if (bagEntry == null || bagEntry.compareAndSet(STATE_NOT_IN_USE, STATE_IN_USE)) {
return bagEntry; // null=超时,由 HikariPool 抛 timeout 异常;拿到还得 CAS,交接的连接也可能被扫描线程截胡
}
timeout -= elapsedNanos(start);
} while (timeout > 10_000); // 剩不到 10 微秒就别等了,直接算超时
return null;
}
finally {
waiters.decrementAndGet();
}
}
public void requite(final T bagEntry)
{
bagEntry.setState(STATE_NOT_IN_USE); // 先无条件置回空闲——从这一刻起任何扫描 sharedList 的线程都能 CAS 抢走它
for (int i = 0; waiters.get() > 0; i++) { // 有人在等就优先手递手,而不是塞回自己的 ThreadLocal
if (bagEntry.getState() != STATE_NOT_IN_USE || handoffQueue.offer(bagEntry)) {
return; // 状态变了=已被扫描线程抢走,或成功交接给一个等待者,都算归还完成
}
else if ((i & 0xff) == 0xff) {
parkNanos(MICROSECONDS.toNanos(10)); // 每 256 次失败睡 10 微秒,让出 CPU 给消费方
}
else {
Thread.yield();
}
}
final List<Object> threadLocalList = threadList.get();
if (threadLocalList.size() < 50) { // 没人等:私藏进本线程缓存,下次 borrow 一级命中;上限 50 防单线程囤积
threadLocalList.add(weakThreadLocals ? new WeakReference<>(bagEntry) : bagEntry);
}
}
public void add(final T bagEntry)
{
if (closed) {
LOGGER.info("ConcurrentBag has been closed, ignoring add()");
throw new IllegalStateException("ConcurrentBag has been closed, ignoring add()");
}
sharedList.add(bagEntry);
// spin until a thread takes it or none are waiting
while (waiters.get() > 0 && bagEntry.getState() == STATE_NOT_IN_USE && !handoffQueue.offer(bagEntry)) {
Thread.yield(); // 新连接也走手递手,确保等待者第一时间被喂到
}
}
// ... remove / reserve / unreserve / values 等房管方法略
}原理串讲
一次完整的借出:业务调 HikariDataSource.getConnection(),双检锁拿到 HikariPool 后进入 HikariPool.getConnection(connectionTimeout),核心就一句 connectionBag.borrow(timeout, MILLISECONDS)。
代码块收起展开
borrow 内部三级查找:先翻本线程的 `threadList`——这是"这个线程上次还回来的连接",纯本地操作,唯一的同步点是一次 `compareAndSet(STATE_NOT_IN_USE, STATE_IN_USE)`;
没命中就顺序扫 `sharedList` 对每个条目试 CAS——注意借出的连接从不离开 sharedList,"借"这个动作没有任何结构性修改,只是把状态位从 0 拨到 1;还没有就 listener.addBagItem 请求池异步补充连接,然后阻塞在 handoffQueue.poll 上等待。
borrow 拿到 PoolEntry 后,HikariPool 还要做存活检查(500ms 内用过的直接跳过),最后 createProxyConnection 包一层代理返回,顺手把泄漏检测任务挂上。
归还是镜像过程:业务调 Connection.close(),代理转发到 HikariPool.recycle → ConcurrentBag.requite,先把状态置回 STATE_NOT_IN_USE,有等待者就通过 handoffQueue 手递手,没有就存进 ThreadLocal 缓存。
为什么”借出不移除”?这是整个设计的支点。传统池借出 = 从空闲队列 remove,归还 = add,每次都是结构性修改,必须加锁或用重量级并发队列。
ConcurrentBag 里连接永远躺在 sharedList,借还只是一个 int 的 CAS,把”容器操作”降维成”状态位翻转”;副产品是房管线程随时能遍历全量连接做 maxLifetime 驱逐、空闲回收(配合 STATE_RESERVED 先占住再操作)。
代价也写在类注释里:借了不还(requite)不会被 GC,必然泄漏——所以泄漏检测只能靠 borrow 时挂的定时任务兜底。
为什么 ThreadLocal 缓存能成立、又为什么 CAS 一次都不能省?连接池的访问模式天然”线程亲和”——一个业务线程反复借还,大概率拿回自己刚还的那个连接,ThreadLocal 让这条最热路径零共享、零竞争。
但 threadList 里存的只是 sharedList 条目的引用副本,别的线程扫 sharedList 时完全可以把你缓存的连接”偷走”(这正是注释里 stolen 的含义),所以即使从自己的缓存里拿,也必须 CAS 验证状态,CAS 失败说明已被偷,继续往下找即可。
偷窃机制反过来保证了公平性下限:缓存只是加速,不是所有权,连接不会被闲置的线程锁死。
为什么归还时优先 handoffQueue 而不是直接进 ThreadLocal?waiters > 0 说明此刻有线程已经扫完 sharedList 正在挨饿,如果归还者把连接私藏进自己的缓存,等待者只能干等到超时。
SynchronousQueue 零容量、fair 模式,offer 成功即意味着某个 poll 中的线程当场拿到了连接,唤醒延迟最小。
requite 里那个自旋 + Thread.yield() + 每 256 次 parkNanos 的节奏,是在”尽快交接”和”别空转烧 CPU”之间的折中。
代码块收起展开
至于"为什么比 dbcp/c3p0 快",ConcurrentBag 的无锁路径是大头,其余是工程细节的堆叠:FastList 替代 ArrayList(去掉 get 的 rangeCheck,remove 从尾部倒着找——正好匹配 JDBC 里 Statement 后开先关的模式);
ProxyConnection 等代理类用 Javassist 在字节码层面生成,方法体极短且 final,利于 JIT 内联;连 if (closed) 这类分支的排布都按 JIT 友好的方向调过。
单点看每项都省不了多少,叠起来就是基准测试里数量级的差距。
设计取舍
- 借出不出容器,状态位代替结构修改:换来无锁与可遍历性,代价是忘还必泄漏,只能靠 leakDetectionThreshold 事后报警。
- sharedList 用 CopyOnWriteArrayList:借还是纯读(遍历+CAS),只有建/驱逐连接才写复制——连接池恰好是”读极多写极少”的场景,换成别的场景这个选择就不成立。
- ThreadLocal 缓存不是所有权:连接可被任意线程偷走,所以缓存命中仍要 CAS;误区是以为一级路径可以不做同步。
- borrow 拿到的只是 PoolEntry”入场券”,存活检查、代理包装、泄漏任务都在 HikariPool 层做——ConcurrentBag 完全不懂 JDBC,它只是个泛型并发容器。
- handoffQueue 选 fair 的 SynchronousQueue:牺牲一点吞吐换等待者 FIFO,防止高并发下个别请求无限饿死。