连接池核心 - 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)

代码块JAVA · 2 行收起展开
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.recycleConcurrentBag.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”之间的折中。

代码块JAVA · 2 行收起展开
至于"为什么比 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,防止高并发下个别请求无限饿死。

延伸阅读