本文目录导读:

我将为您提供一个完整的PHP动态数据源切换方案,包含多种实现方式和最佳实践。
基础动态切换实现
1 使用PDO实现简单切换
<?php
class DynamicDataSource {
private static $instance = null;
private $connections = [];
private $currentConnection = null;
// 数据源配置
private $dataSources = [
'master' => [
'host' => '192.168.1.100',
'port' => 3306,
'database' => 'main_db',
'username' => 'root',
'password' => 'password',
'charset' => 'utf8mb4'
],
'slave' => [
'host' => '192.168.1.101',
'port' => 3306,
'database' => 'main_db',
'username' => 'readonly',
'password' => 'password',
'charset' => 'utf8mb4'
],
'log_db' => [
'host' => '192.168.1.102',
'port' => 3306,
'database' => 'log_db',
'username' => 'log_writer',
'password' => 'password',
'charset' => 'utf8mb4'
]
];
private function __construct() {}
public static function getInstance() {
if (self::$instance === null) {
self::$instance = new self();
}
return self::$instance;
}
/**
* 获取数据库连接
* @param string $dataSourceName
* @param bool $forceNew 是否强制新建连接
* @return PDO
*/
public function getConnection($dataSourceName = 'master', $forceNew = false) {
if (!isset($this->dataSources[$dataSourceName])) {
throw new \InvalidArgumentException("数据源 {$dataSourceName} 不存在");
}
// 如果连接已存在且不需要新建,直接返回
if (!$forceNew && isset($this->connections[$dataSourceName])) {
return $this->connections[$dataSourceName];
}
$config = $this->dataSources[$dataSourceName];
try {
$dsn = sprintf(
"mysql:host=%s;port=%s;dbname=%s;charset=%s",
$config['host'],
$config['port'],
$config['database'],
$config['charset']
);
$conn = new PDO($dsn, $config['username'], $config['password'], [
PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION,
PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC,
PDO::ATTR_EMULATE_PREPARES => false
]);
$this->connections[$dataSourceName] = $conn;
if ($this->currentConnection === null) {
$this->currentConnection = $conn;
}
return $conn;
} catch (PDOException $e) {
// 记录日志
error_log("数据库连接失败: " . $e->getMessage());
throw new \RuntimeException("无法连接数据源 {$dataSourceName}");
}
}
/**
* 动态切换当前连接
* @param string $dataSourceName
* @return PDO
*/
public function switchTo($dataSourceName) {
$conn = $this->getConnection($dataSourceName);
$this->currentConnection = $conn;
return $conn;
}
// 防止克隆
private function __clone() {}
}
2 使用示例
<?php
// 初始化
$dsManager = DynamicDataSource::getInstance();
// 默认使用主库
$conn = $dsManager->getConnection('master');
// 查询用户信息
$stmt = $conn->query("SELECT * FROM users LIMIT 10");
$users = $stmt->fetchAll();
// 切换到从库查询
$readConn = $dsManager->switchTo('slave');
$stmt = $readConn->query("SELECT * FROM orders");
$orders = $stmt->fetchAll();
// 切换到日志库
$logConn = $dsManager->switchTo('log_db');
$logConn->exec("INSERT INTO operation_logs (action, timestamp) VALUES ('test', NOW())");
基于Laravel框架的实现
1 自定义DatabaseManager
<?php
namespace App\Providers;
use Illuminate\Support\ServiceProvider;
use Illuminate\Database\DatabaseManager;
class DynamicDatasourceServiceProvider extends ServiceProvider
{
public function boot()
{
// 注册动态数据库配置
$this->app->singleton('dynamic.db', function ($app) {
return new DynamicDatabaseManager($app);
});
}
public function register()
{
$this->mergeConfigFrom(
__DIR__.'/../config/dynamic_database.php', 'dynamic_database'
);
}
}
class DynamicDatabaseManager
{
protected $app;
protected $connections = [];
protected $currentConnection = null;
public function __construct($app)
{
$this->app = $app;
$this->initDatasources();
}
protected function initDatasources()
{
$datasources = config('dynamic_database.connections');
foreach ($datasources as $name => $config) {
$this->connections[$name] = $config;
}
// 设置默认连接
$this->currentConnection = config('dynamic_database.default');
}
/**
* 切换到指定数据源
* @param string $connection
* @return mixed
*/
public function connection($connection = null)
{
$connection = $connection ?: $this->currentConnection;
if (!isset($this->connections[$connection])) {
throw new \InvalidArgumentException("数据源 {$connection} 不存在");
}
$this->app['config']->set('database.default', $connection);
$this->app['config']->set('database.connections.'.$connection,
$this->connections[$connection]);
// 清除已建立的连接
$this->app['db']->purge($connection);
return $this->app['db']->connection($connection);
}
/**
* 获取所有数据源配置
*/
public function getConnectionConfig()
{
return $this->connections;
}
/**
* 动态添加数据源
*/
public function addConnection($name, $config)
{
$this->connections[$name] = $config;
$this->app['config']->set("database.connections.{$name}", $config);
}
}
2 配置示例 (config/dynamic_database.php)
<?php
return [
'default' => 'master',
'connections' => [
'master' => [
'driver' => 'mysql',
'host' => env('DB_HOST', '127.0.0.1'),
'port' => env('DB_PORT', '3306'),
'database' => env('DB_DATABASE', 'main_db'),
'username' => env('DB_USERNAME', 'root'),
'password' => env('DB_PASSWORD', ''),
'charset' => 'utf8mb4',
'collation' => 'utf8mb4_unicode_ci',
'prefix' => '',
'strict' => true,
],
'slave' => [
'driver' => 'mysql',
'host' => env('SLAVE_DB_HOST', '127.0.0.1'),
'port' => env('SLAVE_DB_PORT', '3306'),
'database' => env('SLAVE_DB_DATABASE', 'main_db'),
'username' => env('SLAVE_DB_USERNAME', 'readonly'),
'password' => env('SLAVE_DB_PASSWORD', ''),
'charset' => 'utf8mb4',
'collation' => 'utf8mb4_unicode_ci',
'prefix' => '',
'strict' => true,
],
'log_db' => [
'driver' => 'mysql',
'host' => env('LOG_DB_HOST', '127.0.0.1'),
'port' => env('LOG_DB_PORT', '3306'),
'database' => env('LOG_DB_DATABASE', 'log_db'),
'username' => env('LOG_DB_USERNAME', 'log_writer'),
'password' => env('LOG_DB_PASSWORD', ''),
'charset' => 'utf8mb4',
'collation' => 'utf8mb4_unicode_ci',
'prefix' => '',
'strict' => true,
]
],
// 读写分离规则
'rule' => [
'select' => ['slave', 'slave'], // 读操作优先从库
'write' => ['master'] // 写操作使用主库
]
];
3 Laravel模型使用示例
<?php
namespace App\Models;
use Illuminate\Database\Eloquent\Model;
class User extends Model
{
protected $table = 'users';
// 默认使用的连接
protected $connection = 'master';
// 可用的数据源
const DATA_SOURCES = ['master', 'slave'];
/**
* 在指定数据源上执行查询
* @param string $connection
* @param \Closure $callback
*/
public static function onConnection($connection, \Closure $callback)
{
$model = new static;
$model->setConnection($connection);
try {
return $callback($model);
} finally {
// 恢复默认连接
$model->setConnection('master');
}
}
/**
* 从从库读取数据
*/
public static function readFromSlave(\Closure $callback)
{
return self::onConnection('slave', $callback);
}
// 使用示例方法
public function getUsersFromSlave()
{
return self::readFromSlave(function ($model) {
return $model->select('*')->limit(10)->get();
});
}
}
高级特性实现
1 基于注解的动态切换
<?php
namespace App\Annotations;
/**
* @Annotation
* @Target({"METHOD"})
*/
class DataSource
{
public $value; // 数据源名称
public function getValue()
{
return $this->value;
}
}
// 使用示例
class UserService
{
/**
* @DataSource("slave")
*/
public function getUsers()
{
// 这个方法会自动使用从库连接
return User::all();
}
/**
* @DataSource("master")
*/
public function createUser()
{
// 这个方法会自动使用主库连接
return User::create($data);
}
}
// 注解解析器
class DatasourceAnnotationHandler
{
public function process($className, $methodName)
{
$reflection = new \ReflectionMethod($className, $methodName);
// 解析注解
$annotations = $this->parseAnnotations($reflection->getDocComment());
if (isset($annotations['DataSource'])) {
$datasource = $annotations['DataSource'];
// 在这里实现动态切换
app('dynamic.db')->connection($datasource);
}
return $annotations;
}
protected function parseAnnotations($docBlock)
{
$annotations = [];
if (preg_match_all('/@(\w+)\("([^"]*)"\)/', $docBlock, $matches)) {
for ($i = 0; $i < count($matches[0]); $i++) {
$annotations[$matches[1][$i]] = $matches[2][$i];
}
}
return $annotations;
}
}
2 自动读写分离
<?php
namespace App\Services;
use Illuminate\Database\Query\Builder;
use Illuminate\Support\Facades\DB;
class AutoReadWriteSplit
{
protected $writeConnection = 'master';
protected $readConnections = ['slave', 'slave2']; // 可以配置多个从库
/**
* 自动选择连接
* @param string $sql
* @param array $bindings
* @return mixed
*/
public function executeQuery($sql, $bindings = [])
{
$sql = trim($sql);
// 判断SQL类型
if ($this->isReadQuery($sql)) {
// 选择从库(可以轮询或随机选择)
$connection = $this->selectReadConnection();
return $this->executeOnConnection($connection, $sql, $bindings);
}
// 写操作使用主库
return $this->executeOnConnection($this->writeConnection, $sql, $bindings);
}
protected function isReadQuery($sql)
{
// 判断是SELECT还是其他
return stripos($sql, 'SELECT') === 0;
}
protected function selectReadConnection()
{
// 轮询选择从库
$connections = $this->readConnections;
$index = array_rand($connections);
return $connections[$index];
}
protected function executeOnConnection($connection, $sql, $bindings)
{
DB::connection($connection)->statement($sql, $bindings);
}
}
// 使用示例
$splitService = new AutoReadWriteSplit();
$users = $splitService->executeQuery('SELECT * FROM users WHERE age > ?', [18]);
连接池管理
<?php
class ConnectionPool
{
private $pool = [];
private $config;
private $maxConnections;
public function __construct($config, $maxConnections = 10)
{
$this->config = $config;
$this->maxConnections = $maxConnections;
}
/**
* 获取连接
*/
public function getConnection()
{
// 如果有空闲连接,直接使用
if (!empty($this->pool)) {
return array_pop($this->pool);
}
// 检查连接数是否超限
if (count($this->pool) >= $this->maxConnections) {
throw new \RuntimeException('连接池已满');
}
// 创建新连接
return $this->createConnection();
}
/**
* 释放连接回连接池
*/
public function releaseConnection($connection)
{
if (count($this->pool) < $this->maxConnections) {
$this->pool[] = $connection;
} else {
// 如果连接池已满,关闭连接
$connection = null;
}
}
protected function createConnection()
{
$dsn = "mysql:host={$this->config['host']};dbname={$this->config['database']}";
return new PDO($dsn, $this->config['username'], $this->config['password'], [
PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION,
PDO::ATTR_PERSISTENT => true // 持久连接
]);
}
/**
* 关闭所有连接
*/
public function closeAll()
{
$this->pool = [];
}
}
最佳实践建议
1 配置中心化
<?php
// 使用环境变量或配置中心管理数据源
return [
'connections' => [
'master' => [
'host' => env('MASTER_DB_HOST'),
'database' => env('MASTER_DB_NAME'),
'username' => env('MASTER_DB_USER'),
'password' => env('MASTER_DB_PASS'),
],
// 其他数据源类似
]
];
2 错误处理与重试机制
<?php
trait DatabaseRetryAware
{
protected $maxRetries = 3;
protected $retryDelay = 100; // 毫秒
protected function executeWithRetry($callback)
{
$attempts = 0;
while ($attempts < $this->maxRetries) {
try {
return $callback();
} catch (\PDOException $e) {
$attempts++;
if ($attempts >= $this->maxRetries) {
throw $e;
}
usleep($this->retryDelay * 1000);
}
}
}
public function query($sql, $params = [])
{
return $this->executeWithRetry(function () use ($sql, $params) {
return $this->connection->query($sql, $params);
});
}
}
3 监控与日志
<?php
class DatasourceMonitor
{
private $stats = [];
public function recordRequest($datasource, $duration, $success)
{
$this->stats[] = [
'datasource' => $datasource,
'duration' => $duration,
'success' => $success,
'timestamp' => time()
];
}
public function getStats($datasource = null)
{
if ($datasource) {
return array_filter($this->stats, function ($item) use ($datasource) {
return $item['datasource'] === $datasource;
});
}
return $this->stats;
}
public function getAverageLatency($datasource)
{
$requests = $this->getStats($datasource);
if (empty($requests)) {
return 0;
}
$totalDuration = array_sum(array_column($requests, 'duration'));
return $totalDuration / count($requests);
}
}
完整使用示例
<?php
// 初始化管理器
$datasourceManager = DynamicDataSource::getInstance();
// 执行主库操作
$conn = $datasourceManager->switchTo('master');
$result = $conn->exec("INSERT INTO orders (user_id, total_amount) VALUES (1, 99.99)");
// 从从库读取数据(读写分离)
$slaveConn = $datasourceManager->switchTo('slave');
$orders = $slaveConn->query("SELECT * FROM orders ORDER BY created_at DESC LIMIT 10")->fetchAll();
// 自动选择数据源(基于业务逻辑)
if ($needTransaction) {
$conn = $datasourceManager->switchTo('master');
$conn->beginTransaction();
// 执行事务操作
$conn->commit();
}
// 使用后释放或保持连接
$datasourceManager->releaseConnection('slave');
这个方案提供了:
- 多种实现方式:从简单的PDO到框架集成
- 读写分离支持:自动路由查询和写入
- 连接池管理:优化资源利用
- 注解驱动:简洁的配置方式
- 监控与日志:便于性能分析
- 错误处理:包含重试机制
选择哪种实现取决于你的具体需求:
- 简单项目:使用PDO实现即可
- Laravel项目:使用ServiceProvider方式
- 高并发场景:需要连接池和读写分离
- 微服务架构:需要考虑配置中心化和监控