本文目录导读:

我来详细介绍ReactPHP在PHP项目中的事件循环机制及应用。
ReactPHP 事件循环基础
什么是ReactPHP
ReactPHP是PHP的事件驱动、非阻塞I/O框架,类似于Node.js的事件循环机制。
安装
composer require react/event-loop
事件循环核心组件
基本事件循环示例
<?php
require 'vendor/autoload.php';
use React\EventLoop\Factory;
use React\EventLoop\Loop;
// 创建事件循环
$loop = Loop::get();
// 添加定时器
$loop->addTimer(2, function () {
echo "2秒后执行\n";
});
// 周期性定时器
$loop->addPeriodicTimer(1, function () {
echo "每秒执行一次\n";
});
// 立即执行
$loop->futureTick(function () {
echo "下一次tick执行\n";
});
// 启动事件循环
$loop->run();
实际应用场景
HTTP服务器
<?php
require 'vendor/autoload.php';
use React\Http\Server;
use React\Http\Response;
use Psr\Http\Message\ServerRequestInterface;
use React\EventLoop\Factory;
$loop = \React\EventLoop\Loop::get();
$server = new Server(function (ServerRequestInterface $request) {
return new Response(
200,
['Content-Type' => 'text/plain'],
"Hello World\n"
);
});
$socket = new React\Socket\Server('127.0.0.1:8080', $loop);
$server->listen($socket);
echo "Server running on http://127.0.0.1:8080\n";
$loop->run();
异步文件处理
<?php
require 'vendor/autoload.php';
use React\EventLoop\Loop;
use React\Filesystem\Filesystem;
$loop = Loop::get();
$filesystem = Filesystem::create($loop);
// 异步读取文件
$filesystem->file('data.txt')->getContents()->then(
function ($contents) {
echo "文件内容: " . $contents . "\n";
},
function ($error) {
echo "读取失败: " . $error->getMessage() . "\n";
}
);
// 异步写入文件
$filesystem->file('output.txt')->putContents('异步写入的内容')->then(
function () {
echo "文件写入成功\n";
}
);
$loop->run();
数据库查询示例(模拟)
<?php
require 'vendor/autoload.php';
use React\EventLoop\Loop;
use React\Promise\Deferred;
class AsyncDatabase {
private $loop;
public function __construct($loop) {
$this->loop = $loop;
}
public function query($sql) {
$deferred = new Deferred();
// 模拟异步数据库查询
$this->loop->addTimer(0.1, function () use ($deferred, $sql) {
// 模拟查询结果
$result = [
'id' => 1,
'name' => 'User ' . rand(1, 100),
'query' => $sql
];
$deferred->resolve($result);
});
return $deferred->promise();
}
}
$loop = Loop::get();
$db = new AsyncDatabase($loop);
// 并行执行多个查询
$promises = [
$db->query("SELECT * FROM users"),
$db->query("SELECT * FROM posts"),
$db->query("SELECT * FROM comments")
];
\React\Promise\all($promises)->then(function ($results) {
echo "所有查询完成:\n";
foreach ($results as $result) {
echo "- " . $result['name'] . "\n";
}
});
$loop->run();
WebSocket服务器
<?php
require 'vendor/autoload.php';
use React\EventLoop\Loop;
use React\Socket\Server;
use React\Socket\ConnectionInterface;
use Ratchet\Server\IoServer;
use Ratchet\Http\HttpServer;
use Ratchet\WebSocket\WsServer;
use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface as RatchetConnection;
class Chat implements MessageComponentInterface {
protected $clients;
public function __construct() {
$this->clients = new \SplObjectStorage;
}
public function onOpen(RatchetConnection $conn) {
$this->clients->attach($conn);
echo "新连接: {$conn->resourceId}\n";
}
public function onMessage(RatchetConnection $from, $msg) {
foreach ($this->clients as $client) {
if ($from !== $client) {
$client->send($msg);
}
}
}
public function onClose(RatchetConnection $conn) {
$this->clients->detach($conn);
echo "连接关闭: {$conn->resourceId}\n";
}
public function onError(RatchetConnection $conn, \Exception $e) {
echo "错误: {$e->getMessage()}\n";
$conn->close();
}
}
$loop = Loop::get();
$webSock = new Server('0.0.0.0:8080', $loop);
$server = new IoServer(
new HttpServer(
new WsServer(
new Chat()
)
),
$webSock
);
echo "WebSocket服务器运行在 ws://0.0.0.0:8080\n";
$loop->run();
高级特性
流式处理
<?php
require 'vendor/autoload.php';
use React\EventLoop\Loop;
use React\Stream\ReadableResourceStream;
use React\Stream\WritableResourceStream;
$loop = Loop::get();
// 读取文件流
$readStream = new ReadableResourceStream(
fopen('large-file.txt', 'r'),
$loop
);
// 写入输出流
$writeStream = new WritableResourceStream(
fopen('php://stdout', 'w'),
$loop
);
// 管道传输
$readStream->pipe($writeStream);
$readStream->on('end', function () {
echo "\n文件读取完成\n";
Loop::get()->stop();
});
$loop->run();
自定义事件
<?php
require 'vendor/autoload.php';
use React\EventLoop\Loop;
use Evenement\EventEmitter;
class DataProcessor extends EventEmitter {
private $loop;
public function __construct($loop) {
$this->loop = $loop;
}
public function process($data) {
$this->emit('processing', [$data]);
$this->loop->addTimer(1, function () use ($data) {
$result = strtoupper($data);
$this->emit('processed', [$result]);
$this->loop->addTimer(0.5, function () use ($result) {
$this->emit('completed', [$result]);
});
});
}
}
$loop = Loop::get();
$processor = new DataProcessor($loop);
$processor->on('processing', function($data) {
echo "处理中: $data\n";
});
$processor->on('processed', function($result) {
echo "处理完成: $result\n";
});
$processor->on('completed', function($result) {
echo "完成: $result\n";
$loop->stop();
});
$processor->process("hello world");
$loop->run();
性能优化建议
定时器管理
<?php
// 使用弱引用避免内存泄漏
$timer = $loop->addPeriodicTimer(1, function () use (&$timer) {
static $count = 0;
$count++;
if ($count >= 10) {
$loop->cancelTimer($timer); // 取消定时器
}
});
错误处理
<?php
// 全局错误处理
$loop->addSignal(SIGINT, function () {
echo "收到退出信号,清理资源...\n";
// 清理操作
exit(0);
});
最佳实践
避免阻塞操作
// ❌ 错误做法
$loop->addPeriodicTimer(1, function () {
sleep(5); // 阻塞事件循环
});
// ✅ 正确做法
$loop->addPeriodicTimer(1, function () use ($loop) {
$loop->addTimer(5, function () {
// 异步执行
});
});
资源管理
<?php
class ResourceMonitor {
private $loop;
private $connections = [];
private $maxConnections = 1000;
public function __construct($loop) {
$this->loop = $loop;
}
public function addConnection($conn) {
if (count($this->connections) >= $this->maxConnections) {
throw new \RuntimeException("达到最大连接数");
}
$this->connections[spl_object_hash($conn)] = $conn;
$conn->on('close', function () use ($conn) {
unset($this->connections[spl_object_hash($conn)]);
});
}
}
测试示例
<?php
use PHPUnit\Framework\TestCase;
use React\EventLoop\Loop;
class EventLoopTest extends TestCase {
public function testTimer() {
$loop = Loop::get();
$executed = false;
$loop->addTimer(0.01, function () use (&$executed) {
$executed = true;
});
$loop->run();
$this->assertTrue($executed);
}
}
ReactPHP的事件循环机制非常适合处理高并发、实时通信、流式数据处理等场景,合理使用可以显著提升PHP应用的性能和响应能力。