Skip to content

Repository files navigation

Universal Connection Pool

一个全能连接池的设计与开发。它不仅是一个 borrow/return 的资源缓存,而是把固定/弹性容量、多 endpoint 路由(读写分离 / 权重 / 一致性哈希 / latency-aware)、健康检查与生命周期治理、多租户配额与优先级、熔断与故障转移、会话/事务/流多路复用、凭证轮换、动态拓扑更新、生产级可观测性与容量顾问,以及可选的 C ABI 和 C++20 协程支持整合到一套统一 API 之中。

对于普通「借一个连接、用完还回去」的场景走最短路径;数据库读写分离、GPU/设备独占、gRPC 流复用、多租户限流、灰度熔断等高级需求按需启用,互不拖累。

  • 语言标准:C++20 必需(协程 awaitable 与 std::stop_token 随之启用)。
  • 依赖:仅标准库 + 线程库。无第三方依赖。
  • 泛型:connection_pool<T> 复用任意长生命周期资源——数据库连接、TCP/HTTP/WebSocket/RPC 客户端、Redis/MQ 客户端、对象存储 SDK、GPU/设备句柄等。

目录


整体架构

采用「facade + pool_state + router + factory + reaper + observability」的分层设计。对外是稳定的借还 API,对内是可替换的路由策略、健康策略与资源契约。

┌─────────────────────────────────────────────────────────────────────┐
│  应用代码                                                              │
│   connection_pool<T> / borrowed_connection<T> / stream_lease /        │
│   transaction_scope / session_guard / shadow_connection /             │
│   capacity_advisor / C ABI (ucp_*)                                    │
└───────────────────────────────┬───────────────────────────────────────┘
                                │  borrow / try_borrow / borrow_for / borrow_async
                                │  borrow_stream / shadow_borrow / begin_transaction
                                ▼
┌─────────────────────────────────────────────────────────────────────┐
│  pool_state<T>(核心,单 mutex_ + cv_ 保护)                          │
│                                                                       │
│  借用层 ── 解析 borrow_options,走 overload 策略与 waiter 授权         │
│     │                                                                 │
│     ▼                                                                 │
│  路由 / 匹配层  endpoint_router                                        │
│   ├── single / round_robin / weighted / read_write_split              │
│   ├── consistent_hash(shard/tenant/affinity key)/ latency-aware      │
│   ├── group/region/tag/capability/device/numa 过滤 + fallback chain   │
│   └── 熔断 open 时按 recovery probe 放量探活                           │
│     │                                                                 │
│     ▼                                                                 │
│  连接容器                                                             │
│   ├── idle_ 空闲队列 + idle_by_endpoint_ 按 endpoint 索引(O(1) 命中) │
│   ├── active_records_ 活跃连接(lease_generation 校验借用句柄)        │
│   ├── waiters_ 等待队列(FIFO / 优先级 / 公平,按 lane 授权)          │
│   └── endpoint_runtime_ / tenant_runtime_ 运行时计数与熔断状态         │
│     │                                                                 │
│     ▼                                                                 │
│  后台 reaper 线程                                                     │
│   ├── 空闲淘汰 / 超龄退休 / 泄漏检测 / 健康巡检                        │
│   ├── endpoint 元数据/凭证过期驱逐                                     │
│   └── 慢 create orphan 收敛、close_timeout 保护                        │
│                                                                       │
│  可观测层:原子计数 + 等待延迟直方图 + endpoint/tenant 快照 +          │
│            capacity_advisor(饱和度 / EWMA 失败率 / 扩缩容建议)       │
└─────────────────────────────────────────────────────────────────────┘

数据流要点:

  • 简单借用先查 idle_,命中即走校验快路径;未命中且未达 max_size 则新建。
  • endpoint 命中借助 idle_by_endpoint_ 索引,指定 endpoint_hint 时按桶查找,避免全量扫描。
  • 等待授权lane_key 分组,同 lane 内按 FIFO 序号或优先级选出授权 waiter,避免惊群。
  • 慢 create:用户 create 回调慢时,连接以 orphan 方式在后台收敛,不阻塞持锁路径;close_timeout 保证关闭可终止。
  • 健康与生命周期由 reaper 线程周期驱动:空闲超时、最大寿命、最大复用、泄漏、endpoint 元数据变化、凭证过期都在此统一清理。
  • 熔断与故障转移在路由层完成:endpoint 连续失败开启熔断,recovery 策略放量探活,fallback_chain / failover group 兜底。

代码结构

include/universal_connection_pool/
├── connection_pool.hpp     伞形头文件 + pool_state<T> 全部核心实现
├── options.hpp             pool_options / borrow_options、枚举、策略契约
├── borrowed_connection.hpp RAII 借用句柄、stream_lease、transaction/session guard、shadow
├── endpoint.hpp            endpoint_config、endpoint_router 路由算法与过滤
├── factory.hpp             connection_factory<T>:create/reset/close/inspect 回调族
├── health_check.hpp        health_checker<T> 健康检查接口
├── advisor.hpp             capacity_advisor:饱和度快照、EWMA 失败率、what-if
├── stats.hpp               pool_stats / endpoint_stats / tenant_stats、直方图
├── async.hpp               异步 / 协程 awaitable 扩展
├── errors.hpp              异常类型(borrow_timeout / pool_closed / invalid_pool_options…)
├── adapters.hpp            适配器聚合头
└── c_api.h                 可选 C ABI,用于 FFI 与 C 调用方

src/
├── connection_pool.cpp     策略模板、apply_overlay、validate_options
├── advisor.cpp             容量顾问实现(饱和度、EWMA、what-if、成本模型)
├── stats.cpp               JSON / Prometheus 序列化
├── c_api.cpp               C ABI 包装(引用计数守卫、opaque handle、边界吞异常)
├── endpoint.cpp / health_check.cpp  header-only 实现的编译单元占位

tests/                      逐特性 CTest:routing / resilience / credential / stream /
                            locality / policy / capacity_advisor / concurrency / c_api …
examples/                   basic / database_like / read_write_split / async_borrow /
                            tcp_client / c_api.c / adapters
benchmarks/connection_pool_bench.cpp  借还吞吐基准
cmake/                      find_package 导出配置

适用场景及原因

场景 推荐用法 为什么适用
通用资源复用(DB / RPC / Redis 客户端) connection_pool<T>::borrow() + 默认固定容量 RAII 自动归还、异常安全关闭,零额外配置即可用
数据库读写分离 routing_policy::read_write_split + read/write 角色 endpoint borrow_intent 路由到读/写节点,权重分摊只读副本
多副本负载均衡 weighted / latency_aware routing 权重票据分摊;latency-aware 按 EWMA 延迟+错误率+在途选最优
分片 / 会话亲和 consistent_hash + shard_key / session_guard pin 同 key 稳定落到同分片;pin 保证会话内复用同一连接
多租户隔离 tenant_quota + priority_class + overflow 策略 每租户配额与突发上限,quota reject 独立计数不污染超时指标
GPU / 设备独占 exclusive endpoint + device_id / numa_node 亲和 独占资源 max=1,按设备/NUMA 亲和路由
gRPC / HTTP2 流复用 borrow_stream + multiplexed 契约 单连接多路复用流,流满自动新建连接
灰度 / 影子流量 shadow_borrow 主连接之外尽力借一个影子连接做双写,超时可配
故障韧性 熔断阈值 + recovery 探活 + fallback_chain / failover 连续失败熔断,放量探活恢复,兜底链路故障转移
凭证轮换(短期 token / IAM) credential_provider + pre-expire drain 旧凭证连接排空并强制退休,不影响在途请求
动态拓扑(服务发现) endpoint_provider + update_endpoints 运行时增删 endpoint,旧连接在 reaper 周期主动清理
过载敏感服务端 overload_policy::fail_fast / custom hook + 指标 显式背压而非无限排队;自定义 overload 决策
弹性容量 min_size / max_size + 动态 resize + capacity_advisor 按需伸缩,顾问基于饱和度与 EWMA 失败率给扩缩容建议
C / 其他语言集成 c_api.hucp_* 稳定 C ABI、opaque handle、append-only struct 兼容、异常不跨边界
C++20 协程 co_await pool.borrow_async() 把借用接入协程,无需手写回调

什么时候适合

  • 无状态、创建极廉价的资源:如果连接创建本身就是纳秒级,池化的簿记开销可能得不偿失。
  • 单次、一过性使用:程序里只借一次的资源,直接构造更省事;池的价值在于复用与治理。
  • 需要完整分布式连接治理平面(跨进程连接编排、全局限流协调):本库定位是「进程内全能但克制」的连接池,不是分布式中间件。

编译与运行

从项目根目录执行:

cmake -B build -DCMAKE_BUILD_TYPE=Debug
cmake --build build -j$(nproc)
cd build && ctest --output-on-failure
cd ..

常用 CMake 选项

# 示例 + 测试(默认均为 ON)
cmake -B build -DUCP_BUILD_EXAMPLES=ON -DUCP_BUILD_TESTS=ON

# 基准(默认 OFF)
cmake -B build-bench -DUCP_BUILD_BENCHMARKS=ON -DUCP_BUILD_TESTS=OFF
cmake --build build-bench --target ucp_connection_pool_bench

Sanitizer 构建(GCC / Clang)

# AddressSanitizer + UndefinedBehaviorSanitizer
cmake -B build-asan -DCMAKE_BUILD_TYPE=Debug -DUCP_ENABLE_ASAN=ON -DUCP_ENABLE_UBSAN=ON
cmake --build build-asan -j$(nproc)
cd build-asan && ctest --output-on-failure && cd ..

# ThreadSanitizer(验证并发正确性,无数据竞争)
cmake --preset tsan && cmake --build --preset tsan && ctest --preset tsan

CMake Preset

cmake --preset debug     && cmake --build --preset debug     && ctest --preset debug
cmake --preset asan-ubsan && cmake --build --preset asan-ubsan && ctest --preset asan-ubsan
cmake --preset tsan      && cmake --build --preset tsan      && ctest --preset tsan

完整本地验收

cmake --preset debug && cmake --build --preset debug && ctest --preset debug
cmake --preset asan-ubsan && cmake --build --preset asan-ubsan && ctest --preset asan-ubsan
cmake --preset tsan && cmake --build --preset tsan && ctest --preset tsan

cmake -B build-bench -DCMAKE_BUILD_TYPE=Release -DUCP_BUILD_BENCHMARKS=ON -DUCP_BUILD_TESTS=OFF -DUCP_BUILD_EXAMPLES=OFF
cmake --build build-bench -j$(nproc)
./build-bench/ucp_connection_pool_bench

cmake --install build --prefix /tmp/ucp-install
cmake -B /tmp/ucp-consumer-build tests/install_smoke -DCMAKE_PREFIX_PATH=/tmp/ucp-install
cmake --build /tmp/ucp-consumer-build -j$(nproc)
/tmp/ucp-consumer-build/ucp_install_smoke

作为依赖被其他 CMake 工程消费

cmake --install build --prefix /opt/universal_connection_pool
find_package(universal_connection_pool CONFIG REQUIRED)
target_link_libraries(app PRIVATE universal_connection_pool::universal_connection_pool)

快速上手示例

最小 C++ 用法:

#include "universal_connection_pool/connection_pool.hpp"
#include <iostream>

struct client {
    int id = 0;
    bool open = true;
    void close() { open = false; }
};

int main() {
    int next_id = 1;

    ucp::connection_factory<client> factory;
    factory.create = [&] { return client{next_id++, true}; };
    factory.close  = [](client& c) { c.close(); };

    ucp::pool_options options;
    options.min_size     = 1;
    options.max_size     = 8;
    options.initial_size = 1;

    ucp::connection_pool<client> pool(factory, options);

    auto borrowed = pool.borrow();          // RAII,作用域结束自动归还
    std::cout << "client id = " << borrowed->id << "\n";
    std::cout << pool.dump_json() << "\n";
}

读写分离路由:

ucp::pool_options options;
options.routing = ucp::routing_policy::read_write_split;
options.endpoints = {
    {.id = "primary", .role = ucp::endpoint_role::write},
    {.id = "replica", .role = ucp::endpoint_role::read, .weight = 2},
};

ucp::borrow_options read;
read.intent = ucp::borrow_intent::read;      // 路由到只读副本
auto conn = pool.borrow(read);

带超时 / 异步借用:

if (auto conn = pool.borrow_for(std::chrono::milliseconds(200))) {
    use(**conn);
}                                            // 超时返回空 optional

auto fut = pool.borrow_async();              // C++20:co_await pool.borrow_async()
auto conn = fut.get();

容量顾问:

ucp::capacity_advisor advisor;
auto snap = advisor.snapshot(pool.stats(), pool.options());
if (snap.saturated) {
    auto rec = advisor.recommend(pool.stats(), pool.options());
    // rec.action / rec.recommended_max_size / rec.reason
}

更多示例(异步取消、熔断、故障转移、多租户、流复用、事务/会话、凭证轮换、C API)见 examples/tests/

C API 示例

#include "universal_connection_pool/c_api.h"
#include <stdio.h>
#include <stdlib.h>

static void* create_conn(void* user_data) {
    int* next_id = (int*)user_data;
    int* value = (int*)malloc(sizeof(int));
    if (value) *value = (*next_id)++;
    return value;
}
static void close_conn(void* connection, void* user_data) {
    (void)user_data;
    free(connection);
}

int main(void) {
    int next_id = 1;
    ucp_pool_options options;
    ucp_pool_options_init(&options);
    options.user_data = &next_id;
    options.create    = create_conn;
    options.close     = close_conn;
    options.max_size  = 4;

    ucp_pool* pool = ucp_pool_create(&options);
    if (!pool) return 1;

    ucp_connection* conn = NULL;
    if (ucp_borrow(pool, &conn, 100) != UCP_OK) {
        ucp_pool_destroy(pool);
        return 2;
    }
    int* value = (int*)ucp_connection_get(conn);
    printf("connection id: %d\n", value ? *value : -1);

    ucp_return(pool, conn);
    ucp_pool_destroy(pool);
    return 0;
}

运行示例

./build/ucp_basic_example
./build/ucp_database_like_example
./build/ucp_read_write_split_example
./build/ucp_async_borrow_example
./build/ucp_tcp_client_example
./build/ucp_c_api_example

# 启用 adapter 示例
cmake -B build -DUCP_BUILD_ADAPTERS=ON && cmake --build build -j$(nproc)
./build/ucp_adapters_example

可观测性

auto s = pool.stats();
s.borrowed_total;               // 借出总数
s.borrow_timeout_total;         // 借用超时(不含租户配额拒绝)
s.tenant_quota_reject_total;    // 租户配额拒绝(独立计数)
s.idle_enqueue_failure_total;   // OOM 等归还入队失败
s.average_wait_time_ns();       // 平均排队等待
for (auto& e : s.endpoints) { /* per-endpoint 计数、熔断、EWMA 延迟 */ }
for (auto& t : s.tenants)   { /* per-tenant 配额与拒绝 */ }

auto json = ucp::dump_json(s);            // JSON
auto prom = ucp::to_prometheus(s, "svc"); // Prometheus 文本
auto diag = pool.dump_diagnostics_json(); // 含策略/拓扑/凭证代次的诊断快照

连接生命周期钩子:on_borrow_queued / on_borrow_acquired / on_connection_returned / on_connection_invalidated / on_health_failed / on_leak_detected / on_sampled_log,可接入 tracing 或自定义监控。

容量顾问 capacity_advisor 输出饱和度快照(利用率、等待压力、超时率、EWMA 创建失败率、endpoint 倾斜)、扩缩容建议与 what-if 成本估算。


许可说明

知识星球:“奔跑中cpp / c++” 所有,

阿甘微信:LLqueww

商业使用前请联系我方授权。一旦发现侵权行为,将依法追究法律责任。

(对于公司法律事务已有对接律师,敬请告知)

About

All-purpose connection pool: fixed/elastic capacity, multi-endpoint routing, health management, tenant quota, circuit breaker, multiplexing, credential rotation, observability, C ABI & C++20 coroutine support.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages