workerman、GatewayWorker 中使用连接池类 ,应用场景高并发处理
·
<?php
namespace app\socket;
use Predis\Client;
use Predis\Connection\ConnectionException;
use think\facade\Db;
use think\facade\Config;
use think\facade\Log;
use Exception;
/**
* 数据库连接池管理类
*
* 功能:
* 1. 管理MySQL和Redis连接的创建、复用和回收
* 2. 动态维护连接池大小(自动扩缩容)
* 3. 定期检测空闲连接健康状态, 在使用类中创建定时任务维护连接池
* 4. 实现LRU策略分配连接
*
* 特性:
* - 支持最小/最大连接数配置
* - 空闲超时自动回收机制
* - 异常连接自动剔除
* - 连接资源争用时的等待策略
*/
class ConnectionPool
{
// 配置标识常量
const MYSQL_CONFIG_NAME = 'workerman_mysql'; // MySQL配置标识
const REDIS_CONFIG_NAME = 'workerman_redis'; // Redis配置标识
const IDLE_THRESHOLD = 5; // 空闲连接复用检测阈值(秒)
// 连接池存储结构
protected $mysqlConnections = []; // MySQL连接池 [连接ID => [connection, lastUsedTime, isInUse]]
protected $redisConnections = []; // Redis连接池 [连接ID => [connection, lastUsedTime, isInUse]]
// 连接池配置
protected $mysqlPoolConfig = []; // MySQL连接池配置(min/max/idle_timeout等)
protected $redisPoolConfig = []; // Redis连接池配置
protected $redisConfig = []; // Redis原生连接参数
/**
* 初始化连接池
*
* 流程:
* 1. 加载配置文件
* 2. 创建最小数量的MySQL连接
* 3. 创建最小数量的Redis连接
*/
public function initialize()
{
$this->loadConfig();
// 初始化MySQL连接池
for ($i = 0; $i < $this->mysqlPoolConfig['min_connections']; $i++) {
$this->createMysqlConnection();
}
// 初始化Redis连接池
for ($i = 0; $i < $this->redisPoolConfig['min_connections']; $i++) {
$this->createRedisConnection();
}
}
/**
* 获取可用MySQL连接
*
* @return \think\db\Connection MySQL连接对象
* @throws Exception 当连接池满载且无可用连接时抛出异常
*/
public function getDbConnection()
{
return $this->acquireConnection(
$this->mysqlConnections,
$this->mysqlPoolConfig,
[$this, 'createMysqlConnection'] // 连接创建回调
);
}
/**
* 获取可用Redis连接
*
* @return \Predis\Client Redis连接对象
* @throws Exception 当连接池满载且无可用连接时抛出异常
*/
public function getRedisConnection()
{
return $this->acquireConnection(
$this->redisConnections,
$this->redisPoolConfig,
[$this, 'createRedisConnection'] // 连接创建回调
);
}
/**
* 加载数据库连接池配置
*
* 分离配置项:
* - 提取MySQL连接池参数
* - 分离Redis连接参数与连接池参数
*/
protected function loadConfig()
{
// 加载MySQL配置
$mysqlConfig = Config::get('database.connections.' . self::MYSQL_CONFIG_NAME, []);
$this->mysqlPoolConfig = $mysqlConfig['pool'];
// 加载Redis配置
$redisFullConfig = Config::get('database.connections.' . self::REDIS_CONFIG_NAME, []);
$this->redisConfig = $redisFullConfig; // 原始连接配置
$this->redisPoolConfig = $redisFullConfig['pool']; // 连接池参数
unset($this->redisConfig['pool']); // 移除pool字段保留纯连接配置
}
/**
* 创建MySQL连接并加入连接池
*
* @return \think\db\Connection 新建的MySQL连接
*/
protected function createMysqlConnection()
{
// 创建独立连接(true表示不共享连接)
$connection = Db::connect(self::MYSQL_CONFIG_NAME, true);
$connectionId = spl_object_id($connection);
// 记录连接状态
$this->mysqlConnections[$connectionId] = [
'connection' => $connection, // 连接对象
'lastUsedTime' => time(), // 最后使用时间戳
'isInUse' => false // 使用状态标记
];
return $connection;
}
/**
* 创建Redis连接并加入连接池
*
* @return \Predis\Client 新建的Redis连接
*/
protected function createRedisConnection()
{
// 创建Predis客户端(开启异常抛出)
$connection = new Client($this->redisConfig, ['exceptions' => true]);
$connection->connect(); // 显式建立连接
$connectionId = spl_object_id($connection);
$this->redisConnections[$connectionId] = [
'connection' => $connection,
'lastUsedTime' => time(),
'isInUse' => false
];
return $connection;
}
/**
* 维护连接池健康状态
*
* 定期执行(如每分钟):
* 1. 检查空闲连接活性
* 2. 回收超时空闲连接
* 3. 补充最小连接数
*/
public function maintainConnectionPools()
{
// MySQL连接池维护
$this->maintainConnectionPool(
$this->mysqlConnections,
$this->mysqlPoolConfig,
// MySQL活性检测回调
function ($conn) {
$conn->query('SELECT 1'); // 发送心跳查询
},
// MySQL关闭回调
function ($conn) {
$conn->close(); // 关闭数据库连接
}
);
// Redis连接池维护
$this->maintainConnectionPool(
$this->redisConnections,
$this->redisPoolConfig,
// Redis活性检测回调
function ($conn) {
// 检测连接状态及PING响应
if (!$conn->isConnected() || $conn->ping() !== true) {
throw new ConnectionException($conn, 'Redis连接异常');
}
},
// Redis关闭回调
function ($conn) {
$conn->disconnect(); // 断开Redis连接
}
);
}
/**
* 连接池维护核心逻辑
*
* @param array &$connectionPool 连接池引用
* @param array $poolConfig 连接池配置
* @param callable $checkCallback 连接活性检测回调
* @param callable $closeCallback 连接关闭回调
*/
protected function maintainConnectionPool(
array &$connectionPool,
array $poolConfig,
callable $checkCallback,
callable $closeCallback
) {
$now = time();
$invalidIds = []; // 待移除连接ID集合
// 计算当前连接池使用率
$currentCount = count($connectionPool);
$usageRatio = $currentCount / max(1, $poolConfig['max_connections']);
// 动态调整空闲连接回收时间
$idleTimeout = $usageRatio < 0.4
? $poolConfig['idle_timeout']
: (int)($usageRatio * ($poolConfig['idle_timeout'] * 4));
// 设置最大超时限制(基础超时的5倍)
$idleTimeout = min($idleTimeout, $poolConfig['idle_timeout'] * 5);
// 记录调试信息(可选)
Log::debug(sprintf(
"连接池维护: 总数=%d, 使用率=%.2f%%, 动态超时=%ds",
$currentCount,
$usageRatio * 100,
$idleTimeout
));
// 第一轮遍历:检测连接状态并标记失效连接
foreach ($connectionPool as $id => $item) {
// 跳过正在使用的连接
if ($item['isInUse']) continue;
try {
// 执行活性检测(可能抛出异常)
$checkCallback($item['connection']);
// 判断是否达到回收条件(超时且大于最小连接数)
if (
$now - $item['lastUsedTime'] > $idleTimeout &&
$currentCount > $poolConfig['min_connections']
) {
$invalidIds[] = $id; // 标记为待回收
}
} catch (Exception $e) {
// 记录异常并标记失效连接
Log::error("连接异常: {$e->getMessage()}");
$invalidIds[] = $id;
}
}
// 第二轮遍历:关闭并移除失效连接
foreach ($invalidIds as $id) {
try {
$closeCallback($connectionPool[$id]['connection']);
} catch (Exception $e) {
Log::error("关闭连接失败: {$e->getMessage()}");
}
unset($connectionPool[$id]);
}
// 第三阶段:补充最小连接数
$isMysqlPool = ($connectionPool === $this->mysqlConnections);
$currentCount = count($connectionPool);
$minRequired = $poolConfig['min_connections'];
while ($currentCount < $minRequired) {
$isMysqlPool ? $this->createMysqlConnection() : $this->createRedisConnection();
$currentCount++;
}
}
/**
* 连接分配核心逻辑(LRU策略)
*
* @param array &$connectionPool 连接池引用
* @param array $poolConfig 连接池配置
* @param callable $createCallback 连接创建方法
* @return object 可用连接对象
* @throws Exception 无可用连接时抛出
*/
protected function acquireConnection(
array &$connectionPool,
array $poolConfig,
callable $createCallback
) {
$now = time();
$lruId = null; // 最近最久未使用连接ID
$oldestTime = PHP_INT_MAX; // 最旧使用时间戳
// 第一轮遍历:寻找最佳候选连接
foreach ($connectionPool as $id => $item) {
// 跳过已被占用的连接
if ($item['isInUse']) continue;
// 更新LRU连接信息
if ($item['lastUsedTime'] < $oldestTime) {
$oldestTime = $item['lastUsedTime'];
$lruId = $id;
}
// 发现满足空闲阈值的连接立即分配
if ($now - $item['lastUsedTime'] > self::IDLE_THRESHOLD) {
$connectionPool[$id]['lastUsedTime'] = $now;
$connectionPool[$id]['isInUse'] = true;
return $item['connection'];
}
}
// 第二阶段:连接池未满时创建新连接
$currentCount = count($connectionPool);
$maxAllowed = $poolConfig['max_connections'];
if ($currentCount < $maxAllowed) {
Log::info("新建连接 (当前: {$currentCount})");
$newConn = call_user_func($createCallback);
$newId = spl_object_id($newConn);
$connectionPool[$newId]['lastUsedTime'] = $now;
$connectionPool[$newId]['isInUse'] = true;
return $newConn;
}
// 第三阶段:分配LRU连接(保底策略)
if ($lruId !== null) {
$connectionPool[$lruId]['lastUsedTime'] = $now;
$connectionPool[$lruId]['isInUse'] = true;
return $connectionPool[$lruId]['connection'];
}
// 所有连接均被占用时抛出异常
throw new Exception("连接池繁忙,无法获取数据库连接");
}
/**
* 释放连接回连接池
*
* @param object $connection 要释放的连接对象
*/
public function releaseConnection($connection)
{
$id = spl_object_id($connection);
// MySQL连接释放
if (isset($this->mysqlConnections[$id])) {
$this->mysqlConnections[$id]['isInUse'] = false;
$this->mysqlConnections[$id]['lastUsedTime'] = time();
return;
}
// Redis连接释放
if (isset($this->redisConnections[$id])) {
$this->redisConnections[$id]['isInUse'] = false;
$this->redisConnections[$id]['lastUsedTime'] = time();
return;
}
// 未知连接警告
Log::warning("警告: 尝试释放未知连接");
}
/**
* 关闭所有连接池
*
* 用于服务关闭时资源清理
*/
public function closeAllConnectionPools()
{
// 关闭所有MySQL连接
foreach ($this->mysqlConnections as $id => $item) {
try {
$item['connection']->close();
} catch (Exception $e) {
Log::error("MySQL关闭异常: {$e->getMessage()}");
}
unset($this->mysqlConnections[$id]);
}
// 关闭所有Redis连接
foreach ($this->redisConnections as $id => $item) {
try {
$item['connection']->disconnect();
} catch (Exception $e) {
Log::error("Redis关闭异常: {$e->getMessage()}");
}
unset($this->redisConnections[$id]);
}
Log::info("所有连接池已关闭");
}
}
<?php
namespace app\socket;
use GatewayWorker\Lib\Gateway;
use Workerman\Worker;
use app\socket\ConnectionPool;
use app\socket\service\ChatService;
use app\socket\service\UserService;
use Exception;
class Events
{
// 公共检查间隔
const POOL_CHECK_INTERVAL = 60;
protected static $pool;
public static function onWorkerStart(Worker $businessWorker)
{
self::$pool = new ConnectionPool();
self::$pool->initialize();
// 启动连接池维护定时器
\Workerman\Lib\Timer::add(self::POOL_CHECK_INTERVAL, function () {
self::$pool->maintainConnectionPools();
});
}
/**
* 处理连接
*/
public static function onConnect($cid)
{
}
/**
* 处理消息
*/
public static function onMessage($cid, $data)
{
$db = null;
$redis = null;
try {
// 解析消息
$data = json_decode($data, true) ?? [];
$content = $data['content'] ?? [];
// 获取连接资源
$db = self::$pool->getDbConnection();
$redis = self::$pool->getRedisConnection();
// 根据消息类型分发处理
switch ($data['type']) {
case 'bind':
(new UserService($db, $redis))->bindUser($cid, $content);
break;
default:
}
} catch (Exception $e) {
Gateway::sendToClient($cid, json_encode(['code' => 500, 'error' => $e->getMessage()]));
} finally {
// 确保释放连接资源
if ($db) self::$pool->releaseConnection($db);
if ($redis) self::$pool->releaseConnection($redis);
}
}
<?php
return [
// 默认使用的数据库连接配置
'default' => env('DB_DRIVER', 'mysql'),
// 自定义时间查询规则
'time_query_rule' => [],
// 自动写入时间戳字段
// true为自动识别类型 false关闭
// 字符串则明确指定时间字段类型 支持 int timestamp datetime date
'auto_timestamp' => true,
// 时间字段取出后的默认时间格式
'datetime_format' => 'Y-m-d H:i:s',
// 时间字段配置 配置格式:create_time,update_time
'datetime_field' => '',
// 数据库连接配置信息
'connections' => [
'mysql' => [
// 数据库类型
'type' => env('DB_TYPE', 'mysql'),
// 服务器地址
'hostname' => env('DB_HOST', '127.0.0.1'),
// 数据库名
'database' => env('DB_NAME', ''),
// 用户名
'username' => env('DB_USER', 'root'),
// 密码
'password' => env('DB_PASS', ''),
// 端口
'hostport' => env('DB_PORT', '3306'),
// 数据库连接参数
'params' => [],
// 数据库编码
'charset' => env('DB_CHARSET', 'utf8mb4'),
// 数据库表前缀
'prefix' => env('DB_PREFIX', ''),
// 数据库部署方式:0 集中式(单一服务器),1 分布式(主从服务器)
'deploy' => 0,
// 数据库读写是否分离 主从式有效
'rw_separate' => false,
// 读写分离后 主服务器数量
'master_num' => 1,
// 指定从服务器序号
'slave_no' => '',
// 是否严格检查字段是否存在
'fields_strict' => true,
// 是否需要断线重连
'break_reconnect' => false,
// 监听SQL
'trigger_sql' => env('APP_DEBUG', true),
// 开启字段缓存
'fields_cache' => false,
],
// Workerman 持久连接配置
'workerman_mysql' => [
// 数据库类型
'type' => env('DB_TYPE', 'mysql'),
// 服务器地址
'hostname' => env('DB_HOST', '127.0.0.1'),
// 数据库名
'database' => env('DB_NAME', ''),
// 用户名
'username' => env('DB_USER', 'root'),
// 密码
'password' => env('DB_PASS', ''),
// 端口
'hostport' => env('DB_PORT', '3306'),
// 持久连接
'persistent' => true,
// 数据库连接参数
'params' => [
\PDO::ATTR_TIMEOUT => 3, // 3秒查询超时
\PDO::ATTR_PERSISTENT => true, // 持久连接
\PDO::ATTR_ERRMODE => \PDO::ERRMODE_EXCEPTION, // 异常模式
],
// 数据库编码
'charset' => env('DB_CHARSET', 'utf8mb4'),
// 数据库表前缀
'prefix' => env('DB_PREFIX', ''),
// 数据库部署方式:0 集中式(单一服务器),1 分布式(主从服务器)
'deploy' => 0,
// 数据库读写是否分离 主从式有效
'rw_separate' => false,
// 读写分离后 主服务器数量
'master_num' => 1,
// 指定从服务器序号
'slave_no' => '',
// 是否严格检查字段是否存在
'fields_strict' => true,
// 是否需要断线重连
'break_reconnect' => true,
// 监听SQL
'trigger_sql' => false,
// 开启字段缓存
'fields_cache' => false,
// 连接池配置
'pool' => [
'min_connections' => 4, // 最小连接数
'max_connections' => 12, // 最大连接数
'idle_timeout' => 60 // 空闲超时(秒)
]
],
// Workerman 持久连接配置
'workerman_redis' => [
'type' => 'predis', // 使用 Predis 驱动
'scheme' => 'tcp', // 连接协议
'host' => '127.0.0.1',
'port' => 6379,
'password' => '', // 密码(无则留空)
'select' => 0, // 数据库索引
'timeout' => 2, // 超时时间(秒)
'persistent' => true, // 启用长连接复用
// 连接池配置
'pool' => [
'min_connections' => 3, // 最小连接数
'max_connections' => 10, // 最大连接数
'idle_timeout' => 30, // 空闲连接回收时间(秒)
],
],
// 更多的数据库配置信息
],
];
更多推荐




所有评论(0)