Skip to content

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 的核心,它负责:

  1. 注册 I/O 事件监听器
  2. 检查 I/O 事件是否就绪
  3. 调用对应的回调函数
  4. 不断循环直到没有待处理的事件

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: 基于 libev

2. 使用 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 扩展

参考链接