reactphp-x / cycle-async-database
Async Cycle Database driver bridging reactphp-x/mysql-pool and wpjscc/database
Package info
github.com/reactphp-x/cycle-async-database
pkg:composer/reactphp-x/cycle-async-database
dev-master
2026-07-16 10:21 UTC
Requires
- php: >=8.1
- clue/reactphp-sqlite: ^1.7
- react/async: ^4
- reactphp-x/mysql-pool: ^2.0
- wpjscc/database: 2.x-dev
Requires (Dev)
- phpunit/phpunit: ^11.5
This package is not auto-updated.
Last update: 2026-07-17 09:25:25 UTC
README
基于 reactphp-x/mysql-pool 和 wpjscc/database 的异步 Cycle Database 驱动。在 ReactPHP Fiber 环境下,内部通过 React\Async\await() 挂起协程,对外保持 Cycle 风格的同步 API($db->query() / $db->execute() / 查询构建器等)。
架构
Cycle DatabaseInterface (同步 API)
→ AsyncDatabase / AsyncDatabaseManager
→ AsyncMysqlDriver (DriverInterface)
→ React\Async\await($pool->query(...))
→ reactphp-x/mysql-pool → react/mysql
特性
- 异步驱动:内部使用连接池与
React\Async\await(),对外暴露同步风格接口 - 兼容 Cycle Database API:支持
select/insert/update/delete/upsert构建器与DatabaseInterface - 事务:支持
$db->transaction()及begin/commit/rollback - 流式查询:驱动层提供
queryStream(),适合大结果集 - 连接池:可配置最小/最大连接数、等待队列与超时
- SQLite:提供
AsyncSQLiteDriver(基于clue/reactphp-sqlite)
安装
composer require reactphp-x/cycle-async-database
要求:PHP 8.1+,MySQL 5.7+/8.0+(Upsert 使用 MySQL 8.0.19+ 行别名语法)。
快速开始
<?php require __DIR__ . '/vendor/autoload.php'; use Cycle\Database\Config as Config; use ReactphpX\CycleAsyncDatabase\AsyncDatabaseManager; use ReactphpX\CycleAsyncDatabase\AsyncMySQLDriverConfig; use ReactphpX\CycleAsyncDatabase\AsyncTcpConnectionConfig; $dbal = new AsyncDatabaseManager(new Config\DatabaseConfig([ 'default' => 'default', 'databases' => [ 'default' => [ 'driver' => 'mysql', 'prefix' => '', ], ], 'connections' => [ 'mysql' => new AsyncMySQLDriverConfig( connection: new AsyncTcpConnectionConfig( database: getenv('DB_NAME') ?: 'test', host: getenv('DB_HOST') ?: '127.0.0.1', port: (int) (getenv('DB_PORT') ?: 3306), charset: getenv('DB_CHARSET') ?: 'utf8mb4', user: getenv('DB_USER') ?: 'root', password: getenv('DB_PASSWORD') ?: '123456', ), options: [ 'minConnections' => (int) (getenv('DB_POOL_MIN') ?: 1), 'maxConnections' => (int) (getenv('DB_POOL_MAX') ?: 10), 'waitQueue' => (int) (getenv('DB_POOL_QUEUE') ?: 100), 'waitTimeout' => (int) (getenv('DB_POOL_TIMEOUT') ?: 0), ], ), ], ])); $db = $dbal->database('default'); // 简单查询 $stmt = $db->query('SELECT 1 AS one'); var_dump($stmt->fetchAll()); // DDL / DML $db->execute('CREATE TABLE IF NOT EXISTS demo_basic (id INT AUTO_INCREMENT PRIMARY KEY, title VARCHAR(255))'); $affected = $db->execute('INSERT INTO demo_basic (title) VALUES (?)', ['hello']); echo "Inserted rows: {$affected}\n"; echo 'Last ID: ' . $db->getDriver()->lastInsertID() . "\n";
运行示例
composer install export DB_HOST=127.0.0.1 export DB_PORT=3306 export DB_NAME=test export DB_USER=root export DB_PASSWORD=123456 export DB_CHARSET=utf8mb4 php examples/mysql_basic.php php examples/mysql_queries.php php examples/mysql_transactions.php php examples/mysql_upsert.php php examples/sqlite_basic.php
| 变量 | 默认值 | 说明 |
|---|---|---|
DB_HOST |
127.0.0.1 |
MySQL 地址 |
DB_PORT |
3306 |
MySQL 端口 |
DB_NAME |
test |
数据库名 |
DB_USER |
root |
用户名 |
DB_PASSWORD |
123456 |
密码 |
DB_CHARSET |
utf8mb4 |
字符集 |
DB_POOL_MIN |
1 |
连接池最小连接数 |
DB_POOL_MAX |
10 |
连接池最大连接数 |
DB_POOL_QUEUE |
100 |
等待队列容量 |
DB_POOL_TIMEOUT |
0 |
获取连接超时(毫秒,0 表示不超时) |
使用方式
查询构建器
$rows = $db->select('*') ->from('users') ->where(['status' => 'active']) ->orderBy('id', 'DESC') ->limit(10) ->fetchAll(); $db->insert('users')->values(['email' => 'a@b.com', 'name' => 'Adam'])->run(); $db->update('users', ['name' => 'Updated'], ['id' => 1])->run(); $db->delete('users', ['id' => 100])->run();
事务
$result = $db->transaction(function ($txDb) { $txDb->execute('INSERT INTO logs (val) VALUES (?)', [100]); $txDb->execute('INSERT INTO logs (val) VALUES (?)', [200]); return 'ok'; });
UPSERT
wpjscc/database 2.x 要求指定冲突列:
$db->upsert('users') ->columns('email', 'name') ->conflicts('email') ->values(['email' => 'a@b.com', 'name' => 'Adam']) ->run(); $db->upsert('users') ->columns('email', 'name', 'age') ->conflicts('email') ->updates('age', 'name') ->values([ ['email' => 'a@b.com', 'name' => 'Charlie', 'age' => 10], ]) ->run();
流式查询
use ReactphpX\CycleAsyncDatabase\AsyncMysqlDriver; /** @var AsyncMysqlDriver $driver */ $driver = $db->getDriver(); $stream = $driver->queryStream('SELECT * FROM big_table');
核心类
| 类 | 说明 |
|---|---|
AsyncDatabaseManager |
数据库工厂与管理器 |
AsyncDatabase |
实现 DatabaseInterface 的高层封装 |
AsyncMysqlDriver |
基于 mysql-pool 的异步 MySQL 驱动 |
AsyncTransactionDriver |
事务内持有单连接的驱动 |
AsyncStatement |
将 MysqlResult 适配为 StatementInterface |
AsyncMySQLDriverConfig / AsyncTcpConnectionConfig |
MySQL 连接配置 |
AsyncSQLiteDriver / AsyncSQLiteDriverConfig |
SQLite 异步驱动 |
测试
composer install ./vendor/bin/phpunit
License
MIT