PHP 怎么改造阻塞IO

wen PHP项目 2

本文目录导读:

PHP 怎么改造阻塞IO

  1. 使用非阻塞流模式
  2. 使用 stream_select 轮询
  3. 使用 Swoole 扩展(推荐)
  4. 使用 ReactPHP(纯 PHP 实现)
  5. 使用 Amp 库
  6. 数据库 I/O 优化
  7. 文件 I/O 优化
  8. 混合使用异步和同步
  9. 最佳实践建议

在 PHP 中改造阻塞 I/O 是一个常见的优化需求,我将从几个层面来介绍解决方案:

使用非阻塞流模式

设置非阻塞模式

<?php
// 设置流为非阻塞模式
stream_set_blocking($stream, false);
// 或者使用 fcntl
$flags = fcntl($stream, F_GETFL, 0);
fcntl($stream, F_SETFL, $flags | O_NONBLOCK);

示例:非阻塞读取

<?php
$server = stream_socket_server("tcp://0.0.0.0:8080", $errno, $errstr);
stream_set_blocking($server, false);
while (true) {
    $client = @stream_socket_accept($server, 0);
    if ($client) {
        stream_set_blocking($client, false);
        // 处理客户端连接
    }
    // 继续处理其他任务
    usleep(1000);
}

使用 stream_select 轮询

<?php
$server = stream_socket_server("tcp://0.0.0.0:9000", $errno, $errstr);
$clients = [$server];
while (true) {
    $read = $clients;
    $write = null;
    $except = null;
    $ready = stream_select($read, $write, $except, 0, 200000); // 200ms超时
    if ($ready > 0) {
        foreach ($read as $socket) {
            if ($socket === $server) {
                $client = stream_socket_accept($server);
                $clients[] = $client;
            } else {
                $data = fread($socket, 8192);
                if (strlen($data) > 0) {
                    // 处理数据
                }
            }
        }
    }
}

使用 Swoole 扩展(推荐)

异步方式

<?php
// 安装 swoole: pecl install swoole
$server = new Swoole\Server('0.0.0.0', 9501);
$server->on('connect', function ($server, $fd) {
    echo "Client: {$fd} connected\n";
});
$server->on('receive', function ($server, $fd, $reactor_id, $data) {
    // 异步处理
    $server->send($fd, "Received: " . $data);
});
$server->on('close', function ($server, $fd) {
    echo "Client: {$fd} closed\n";
});
$server->start();

协程方式

<?php
Co\run(function () {
    $server = new Co\Socket(AF_INET, SOCK_STREAM, 0);
    $server->bind('0.0.0.0', 9502);
    $server->listen();
    while (true) {
        $client = $server->accept();
        go(function () use ($client) {
            while (true) {
                $data = $client->recv();
                if (!$data) break;
                $client->send("Echo: " . $data);
            }
        });
    }
});

使用 ReactPHP(纯 PHP 实现)

<?php
// 安装: composer require react/socket
require __DIR__ . '/vendor/autoload.php';
use React\EventLoop\Factory;
use React\Socket\ConnectionInterface;
use React\Socket\SocketServer;
$loop = Factory::create();
$server = new SocketServer('0.0.0.0:8080', $loop);
$server->on('connection', function (ConnectionInterface $connection) {
    echo "New connection from: " . $connection->getRemoteAddress() . PHP_EOL;
    $connection->on('data', function ($data) use ($connection) {
        echo "Received: " . $data;
        $connection->write("Echo: " . $data);
    });
    $connection->on('close', function () {
        echo "Connection closed\n";
    });
});
echo "Server running at 0.0.0.0:8080\n";
$loop->run();

使用 Amp 库

<?php
// 安装: composer require amphp/amp amphp/socket
require __DIR__ . '/vendor/autoload.php';
use Amp\Loop;
use Amp\Socket\Socket;
Loop::run(function () {
    $server = Amp\Socket\listen('0.0.0.0:8080');
    while ($socket = yield $server->accept()) {
        Amp\asyncCall(function () use ($socket) {
            while (null !== $chunk = yield $socket->read()) {
                yield $socket->write("Echo: " . $chunk);
            }
        });
    }
});

数据库 I/O 优化

<?php
// 使用 async MySQL 查询
$db = new MysqliDb();
$db->setAsyncMode(true);
$result = $db->queryAsync("SELECT * FROM users WHERE id > ?", [1000]);
// 先处理其他业务
do_other_work();
// 稍后获取结果
if ($result) {
    $rows = $db->fetchAll(true); // 等待并获取结果
    process_users($rows);
}

文件 I/O 优化

<?php
// 使用非阻塞文件读取
$fp = fopen('/path/to/large_file.log', 'r');
stream_set_blocking($fp, false);
while (!feof($fp)) {
    $line = fgets($fp, 4096);
    if ($line !== false) {
        process_line($line);
    } else {
        // 没有数据,先做其他事情
        do_other_work();
        usleep(10000); // 10ms 后重试
    }
}

混合使用异步和同步

<?php
// 处理多个数据库连接并行
class AsyncDBHandler {
    private $connections = [];
    public function asyncQuery($dbConfig, $query) {
        $conn = new PDO();
        $conn->query($query);
        $this->connections[] = $conn;
    }
    public function waitForResults() {
        while ($this->connections) {
            foreach ($this->connections as $key => $conn) {
                if ($conn->inTransaction()) {
                    // 检查结果是否可用
                }
            }
            usleep(1000);
        }
    }
}

最佳实践建议

  1. 根据场景选择方案

    • 简单场景:使用 stream_select 或非阻塞模式
    • 高并发:选择 Swoole 或 Workerman
    • 需要跨平台:使用 ReactPHP 或 Amp
  2. 性能考虑

    // 设置合理的超时时间
    stream_set_timeout($stream, 0, 100000); // 100ms
    // 使用缓冲
    stream_set_write_buffer($stream, 8192);
  3. 错误处理

    // 非阻塞模式下的错误处理
    if (false === fread($stream, 8192)) {
        // 检查是否是 EAGAIN/EWOULDBLOCK
        $err = socket_last_error();
        if ($err !== SOCKET_EAGAIN && $err !== SOCKET_EWOULDBLOCK) {
            // 真正的错误
            throw new Exception("Read error: " . socket_strerror($err));
        }
    }
  4. 监控和调试

    // 使用断言确保非阻塞模式生效
    assert(!stream_get_meta_data($stream)['blocked']);
    // 添加调试日志
    $start = microtime(true);
    if (false !== $data = fread($stream, 8192)) {
        $elapsed = microtime(true) - $start;
        if ($elapsed > 0.1) {
            error_log("Slow read: {$elapsed}s");
        }
    }

非阻塞 I/O 改造的核心是不等待,通过异步处理、事件循环、协程等机制,让程序在等待 I/O 时能够执行其他任务,选择具体方案时要考虑:

  • 项目复杂度
  • 性能需求
  • 团队技术栈
  • 运维成本

如果是新项目,强烈推荐使用 Swoole 或 Workerman,它们提供了更完善的非阻塞 I/O 支持。

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