ReactPHP 事件驱动
概述
ReactPHP 是一个低级的事件驱动库,为 PHP 提供了异步 I/O 能力。它基于事件循环(Event Loop)模型,通过 Promise 和 Stream 实现非阻塞 I/O 操作。与 Swoole 不同,ReactPHP 纯粹使用 PHP 实现,无需特殊的扩展,可以在任何 PHP 环境中运行(包括 Windows),是 PHP 异步编程的标杆项目。
PHP 版本要求
ReactPHP 3.x 支持 PHP 8.1+。本文基于 ReactPHP 3.x 和 PHP 8.1+ 编写。
基础概念
事件循环(Event Loop)
事件循环是 ReactPHP 的核心,它负责:
- 注册 I/O 事件监听器
- 检查 I/O 事件是否就绪
- 调用对应的回调函数
- 不断循环直到没有待处理的事件
Promise
Promise 是异步操作结果的占位符,它有三种状态:
| 状态 | 说明 |
|---|---|
| Pending | 等待中(初始状态) |
| Fulfilled | 已完成(操作成功) |
| Rejected | 已拒绝(操作失败) |
Stream
Stream 是 ReactPHP 中处理数据流的抽象,支持可读流、可写流和双工流。
安装与配置
bash
# 安装 ReactPHP 核心包
composer require react/event-loop
# 安装常用组件
composer require react/http
composer require react/promise
composer require react/socket
composer require react/stream详细说明
EventLoop 事件循环
php
<?php
declare(strict_types=1);
use React\EventLoop\Loop;
// 获取默认事件循环
$loop = Loop::get();
// 添加定时器(一次性)
$loop->addTimer(3.0, function (): void {
echo "3 秒后执行" . PHP_EOL;
});
// 添加定时器(周期性)
$loop->addPeriodicTimer(1.0, function (): void {
echo "Tick: " . date('H:i:s') . PHP_EOL;
});
// 添加信号监听
$loop->addSignal(SIGINT, function (): void {
echo "收到 SIGINT,停止事件循环" . PHP_EOL;
Loop::stop();
});
// 添加未来定时器(在特定时间执行)
$loop->addTimer(5.0, function (): void {
echo "5 秒后执行" . PHP_EOL;
});
echo "事件循环启动..." . PHP_EOL;
$loop->run();
echo "事件循环结束" . PHP_EOL;Promise 异步操作
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\Promise\Promise;
/**
* 创建 Promise
*/
function fetchUser(int $id): Promise
{
return new Promise(function ($resolve, $reject) use ($id): void {
// 模拟异步操作
Loop::addTimer(1.0, function () use ($id, $resolve, $reject): void {
if ($id > 0) {
$resolve(['id' => $id, 'name' => "User {$id}"]);
} else {
$reject(new \InvalidArgumentException('无效的用户 ID'));
}
});
});
}
/**
* Promise 链式调用
*/
fetchUser(1)
->then(
function (array $user): array {
echo "获取用户: {$user['name']}" . PHP_EOL;
return $user;
},
function (\Throwable $error): void {
echo "错误: " . $error->getMessage() . PHP_EOL;
}
);
/**
* Promise 组合
*/
use React\Promise\Deferred;
function all(iterable $promises): Promise
{
$results = [];
$count = count($promises);
$deferred = new Deferred();
if ($count === 0) {
$deferred->resolve([]);
return $deferred->promise();
}
foreach ($promises as $i => $promise) {
$promise->then(
function ($value) use ($i, &$results, $count, $deferred): void {
$results[$i] = $value;
if (count($results) === $count) {
$deferred->resolve($results);
}
},
function (\Throwable $error) use ($deferred): void {
$deferred->reject($error);
}
);
}
return $deferred->promise();
}
// 并发获取多个用户
$promise1 = fetchUser(1);
$promise2 = fetchUser(2);
$promise3 = fetchUser(3);
all([$promise1, $promise2, $promise3])
->then(function (array $users): void {
echo "获取到 " . count($users) . " 个用户" . PHP_EOL;
});
Loop::run();HTTP Client
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\Http\Browser;
use React\Promise\Promise;
$loop = React\EventLoop\Loop::get();
$browser = new Browser($loop);
/**
* GET 请求
*/
$browser->get('https://httpbin.org/get')
->then(
function (Psr\Http\Message\ResponseInterface $response): void {
echo "状态码: " . $response->getStatusCode() . PHP_EOL;
echo "内容: " . $response->getBody()->getContents() . PHP_EOL;
},
function (\Throwable $error): void {
echo "请求失败: " . $error->getMessage() . PHP_EOL;
}
);
/**
* 并发请求
*/
function concurrentRequests(Browser $browser, array $urls): Promise
{
$promises = array_map(
fn(string $url) => $browser->get($url)
->then(
fn(Psr\Http\Message\ResponseInterface $res) => [
'url' => $url,
'status' => $res->getStatusCode(),
],
fn(\Throwable $e) => ['url' => $url, 'error' => $e->getMessage()]
),
$urls
);
return \React\Promise\all($promises);
}
$urls = [
'https://httpbin.org/get?r=1',
'https://httpbin.org/get?r=2',
'https://httpbin.org/get?r=3',
];
concurrentRequests($browser, $urls)
->then(function (array $results): void {
foreach ($results as $result) {
echo $result['url'] . " → " . ($result['status'] ?? 'ERROR') . PHP_EOL;
}
Loop::stop();
});
$loop->run();HTTP Server
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\Http\Message\Response;
use React\Http\Message\ServerRequest;
use React\Http\Server;
use React\Socket\Server as SocketServer;
$loop = React\EventLoop\Loop::get();
$server = new Server($loop, function (ServerRequest $request): Response {
$path = $request->getUri()->getPath();
$method = $request->getMethod();
$body = json_encode([
'path' => $path,
'method' => $method,
'time' => date('Y-m-d H:i:s'),
], JSON_UNESCAPED_UNICODE);
return new Response(
200,
['Content-Type' => 'application/json'],
$body
);
});
$socket = new SocketServer('0.0.0.0:8080', $loop);
$server->listen($socket);
echo "ReactPHP HTTP 服务器运行在 http://0.0.0.0:8080" . PHP_EOL;Socket 连接
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\Socket\ConnectionInterface;
use React\Socket\Server;
$loop = React\EventLoop\Loop::get();
// TCP 服务器
$socket = new Server('0.0.0.0:9000', $loop);
$socket->on('connection', function (ConnectionInterface $connection): void {
echo "新连接" . PHP_EOL;
$connection->on('data', function (string $data) use ($connection): void {
echo "收到: {$data}" . PHP_EOL;
$connection->write("Echo: {$data}");
});
$connection->on('close', function (): void {
echo "连接关闭" . PHP_EOL;
});
$connection->on('error', function (\Throwable $error): void {
echo "错误: " . $error->getMessage() . PHP_EOL;
});
});
echo "TCP 服务器运行在 tcp://0.0.0.0:9000" . PHP_EOL;
$loop->run();Stream 数据流
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\Stream\ReadableResourceStream;
use React\Stream\WritableResourceStream;
$loop = React\EventLoop\Loop::get();
// 从 stdin 读取
$readable = new ReadableResourceStream(STDIN, $loop);
// 写入 stdout
$writable = new WritableResourceStream(STDOUT, $loop);
// 管道连接
$readable->pipe($writable);
// 数据处理管道
$readable->on('data', function (string $chunk): void {
echo "收到 " . strlen($chunk) . " 字节数据" . PHP_EOL;
});
$readable->on('end', function (): void {
echo "输入结束" . PHP_EOL;
Loop::stop();
});
$readable->on('error', function (\Throwable $e): void {
echo "错误: " . $e->getMessage() . PHP_EOL;
});
$loop->run();实战示例
异步任务队列
php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use React\EventLoop\Loop;
use React\Promise\Deferred;
use React\Promise\Promise;
/**
* 异步任务队列
*/
class AsyncTaskQueue
{
private \SplQueue $queue;
private int $concurrency;
private int $running = 0;
public function __construct(int $concurrency = 3)
{
$this->queue = new \SplQueue();
$this->concurrency = $concurrency;
}
/**
* 添加任务
*/
public function enqueue(callable $task): Promise
{
$deferred = new Deferred();
$this->queue->push([$task, $deferred]);
$this->processNext();
return $deferred->promise();
}
private function processNext(): void
{
while ($this->running < $this->concurrency && !$this->queue->isEmpty()) {
$this->running++;
[$task, $deferred] = $this->queue->dequeue();
$task()->then(
function ($result) use ($deferred): void {
$deferred->resolve($result);
},
function (\Throwable $error) use ($deferred): void {
$deferred->reject($error);
},
function () use ($deferred): void {
$this->running--;
$this->processNext();
}
);
}
}
}
/**
* 创建异步任务
*/
function asyncTask(int $id, int $delay): Promise
{
return new Promise(function ($resolve) use ($id, $delay): void {
Loop::addTimer($delay / 1000.0, function () use ($id, $resolve): void {
echo "任务 #{$id} 完成" . PHP_EOL;
$resolve("result_{$id}");
});
});
}
// 使用示例
$queue = new AsyncTaskQueue(3);
$queue->enqueue(fn() => asyncTask(1, 500));
$queue->enqueue(fn() => asyncTask(2, 300));
$queue->enqueue(fn() => asyncTask(3, 800));
$queue->enqueue(fn() => asyncTask(4, 200));
$queue->enqueue(fn() => asyncTask(5, 600));
Loop::addTimer(3, fn() => Loop::stop());
Loop::run();注意事项
阻塞操作
避免阻塞事件循环
在 ReactPHP 的事件循环中,绝不能执行阻塞操作(如 sleep()、file_get_contents()、mysqli_query())。所有 I/O 操作都应使用异步版本或使用 React\ChildProcess 运行外部进程。
内存泄漏
php
<?php
declare(strict_types=1);
// 定时器未清理导致内存泄漏
$timer = $loop->addPeriodicTimer(1.0, function () use (&$timer): void {
// 必须在某个条件下停止
});
// 正确做法:在不需要时清理
// $loop->cancelTimer($timer);错误处理
php
<?php
declare(strict_types=1);
// Promise 必须始终处理 rejected 状态
$promise->then(
fn($result) => handleSuccess($result),
fn(\Throwable $error) => handleError($error) // 不要遗漏
);
// 或使用 done() 自动传播 rejected
$promise->done(
fn($result) => handleSuccess($result)
// rejected 会抛出异常到事件循环
);最佳实践
1. 选择合适的事件循环实现
php
<?php
// 自动选择最佳实现
$loop = React\EventLoop\Loop::get();
// 或手动选择
// - ExtEventLoop: 基于 libevent(性能最佳)
// - StreamSelectLoop: 纯 PHP 实现(兼容性最佳)
// - EvLoop: 基于 libev2. 使用 Zone 管理定时器和流
php
<?php
use React\EventLoop\Loop;
// Zone 可以批量取消定时器
$zone = Loop::get();
$zone->addTimer(1.0, fn() => echo "timer 1\n");
$zone->addTimer(2.0, fn() => echo "timer 2\n");
// 取消 Zone 中所有定时器
// $zone->cancelAll();下一节
继续学习:Parallel 扩展