PHP动态数据源切换方案

wen PHP项目 1

本文目录导读:

PHP动态数据源切换方案

  1. 基础动态切换实现
  2. 基于Laravel框架的实现
  3. 高级特性实现
  4. 连接池管理
  5. 最佳实践建议
  6. 完整使用示例

我将为您提供一个完整的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');

这个方案提供了:

  1. 多种实现方式:从简单的PDO到框架集成
  2. 读写分离支持:自动路由查询和写入
  3. 连接池管理:优化资源利用
  4. 注解驱动:简洁的配置方式
  5. 监控与日志:便于性能分析
  6. 错误处理:包含重试机制

选择哪种实现取决于你的具体需求:

  • 简单项目:使用PDO实现即可
  • Laravel项目:使用ServiceProvider方式
  • 高并发场景:需要连接池和读写分离
  • 微服务架构:需要考虑配置中心化和监控

抱歉,评论功能暂时关闭!