PHP项目ReactPHP与事件循环

wen PHP项目 1

本文目录导读:

PHP项目ReactPHP与事件循环

  1. ReactPHP 事件循环基础
  2. 事件循环核心组件
  3. 实际应用场景
  4. 高级特性
  5. 性能优化建议
  6. 最佳实践
  7. 测试示例

我来详细介绍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应用的性能和响应能力。

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