PHP Fiber 实战:手写一个协程调度器把同步代码跑出并发

2026-09-26 0 699

手头有个小脚本,需要从八个供应商接口拉数据,汇总成一份对账表。每个接口响应大概三五百毫秒,串起来跑完要三秒多。产品那边的要求是”打开页面就能看到”,三秒显然不行。

最先想到的是多线程或者多进程。但 PHP 的多进程要装 pcntl 扩展,多线程要装 parallel 或 pthreads,服务器上未必有。而且进程之间传数据还得序列化,改起来不轻松。

第二个想法是 curl_multi。这个方案对纯 HTTP 是成立的,但脚本里还有一部分是从 Redis 拉配置、读本地文件、算一些统计数据,这些混在一起之后,curl_multi 就力不从心了——它只能管 curl,别的 IO 还是串行。

后来想到 Fiber。这篇文章就复盘一下:用 Fiber 从零搭一个最小的协程调度器,需要多少代码,会遇到什么问题,以及它到底适不适合你的项目。

一、Fiber 是什么,不是什么

先纠正一个常见误解。Fiber 不是并发,它是”可暂停的函数”。

$fiber = new Fiber(function () {
    echo "An";
    Fiber::suspend();
    echo "Bn";
});

echo "1n";
$fiber->start();
echo "2n";
$fiber->resume();
echo "3n";

输出是:

1
A
2
B
3

整个过程中只有一个线程在跑,一个时刻只有一行代码在执行。Fiber 做的事情是:把函数执行到一半的状态(局部变量、调用栈、当前指令位置)保存下来,然后从外层继续往下走;等你想继续的时候,再把现场恢复回去。

状态保存这件事,PHP 引擎在内部做得很彻底。你不用手动 yield 变量,不用把状态拆成数组传递,写法上就跟普通函数一样。

那这跟并发有什么关系?关系在于:程序里大量时间是在等 IO。等数据库返回、等 API 响应、等文件读完。串行执行的时候,这段时间 CPU 是闲着的。如果能在等待的时候先去做别的事,等结果回来了再切回来,整体耗时就能压下来。

Fiber 提供了”切换”的能力,但切换的时机、切换给谁,得由你决定。这个决策者就是调度器。Fiber 是原材料,调度器是机器。

二、第一版:最小的轮转调度器

先写一个能跑起来的东西。目标很简单:多个任务轮流执行,每个任务自己决定什么时候让出控制权。

<?php
declare(strict_types=1);

final class Scheduler
{
    /** @var SplQueue<Fiber> */
    private SplQueue $ready;

    public function __construct()
    {
        $this->ready = new SplQueue();
    }

    public function spawn(Closure $task): void
    {
        $this->ready->enqueue(new Fiber($task));
    }

    public function run(): void
    {
        while (!$this->ready->isEmpty()) {
            /** @var Fiber $fiber */
            $fiber = $this->ready->dequeue();

            if ($fiber->isTerminated()) {
                continue;
            }

            $signal = $fiber->isStarted()
                ? $fiber->resume()
                : $fiber->start();

            // suspend 之后如果任务还没结束,放回队尾
            if ($fiber->isSuspended()) {
                $this->ready->enqueue($fiber);
            }
        }
    }
}

跑一下:

$loop = new Scheduler();

$loop->spawn(function () {
    for ($i = 0; $i < 3; $i++) {
        echo "任务 A: 第 {$i} 步n";
        Fiber::suspend();
    }
});

$loop->spawn(function () {
    for ($i = 0; $i < 3; $i++) {
        echo "任务 B: 第 {$i} 步n";
        Fiber::suspend();
    }
});

$loop->run();

输出会交错出现:A 的第一步、B 的第一步、A 的第二步、B 的第二步……这正是协程的基本形态。

这段代码里有几个点需要说明。

SplQueue 是先进先出的双向链表,dequeue 从头部取,enqueue 从尾部放。用它而不是数组,是因为队列在循环里会频繁进出,数组的 shift 每次都要重排索引,任务一多开销就上来了。

$fiber->start() 和 $fiber->resume() 的区别在于,前者只能调用一次,用于启动一个还没开始的任务;后者用来唤醒一个已经挂起的任务。两个方法的返回值都是”Fiber 内部 suspend() 时传出来的值”,如果任务跑完正常返回了,那返回值就是任务的 return 结果。

这个返回值目前我们没用,但它正是扩展调度器的关键。后面加定时器和 IO 等待的时候,全靠它来传递信号。

三、第二版:加上定时器

第一版的 Fiber::suspend() 不带参数,语义是”我让一下,待会儿再叫我”。但真实场景里更常见的是”我等 200 毫秒之后继续”。

问题来了:当 Fiber 因为等待而挂起时,调度器怎么知道它该放进”立即执行”的队列,还是”等一会儿再说”的队列?

答案是用 suspend() 的参数做标记。

final class Scheduler
{
    private SplQueue $ready;

    /** @var array<string, Fiber> 到期时间(字符串键) => Fiber */
    private array $timers = [];

    public function __construct()
    {
        $this->ready = new SplQueue();
    }

    public function spawn(Closure $task): void
    {
        $this->ready->enqueue(new Fiber($task));
    }

    public function run(): void
    {
        while (true) {
            $this->flushTimers();

            if ($this->ready->isEmpty()) {
                if (!$this->idle()) {
                    return;
                }
                continue;
            }

            $fiber = $this->ready->dequeue();
            if ($fiber->isTerminated()) {
                continue;
            }

            $signal = $fiber->isStarted()
                ? $fiber->resume()
                : $fiber->start();

            // 只有"主动让出"的才立刻回队列
            // 带 timer / io 标记的,已经登记到别处了
            if ($fiber->isSuspended() && $signal === null) {
                $this->ready->enqueue($fiber);
            }
        }
    }

    private function flushTimers(): void
    {
        if (!$this->timers) {
            return;
        }

        $now = microtime(true);
        foreach ($this->timers as $at => $fiber) {
            if ((float) $at <= $now) {
                unset($this->timers[$at]);
                $this->ready->enqueue($fiber);
            }
        }
    }

    /**
     * @return bool 是否还有活要干
     */
    private function idle(): bool
    {
        if (!$this->timers) {
            return false;
        }

        $next = min(array_map('floatval', array_keys($this->timers)));
        $wait = max(0.0, $next - microtime(true));
        usleep((int) ($wait * 1_000_000));

        return true;
    }

    public function sleep(float $seconds): void
    {
        $fiber = Fiber::getCurrent();

        // 不在协程里,退化成真的睡
        if ($fiber === null) {
            usleep((int) ($seconds * 1_000_000));
            return;
        }

        $key = sprintf('%.6f', microtime(true) + $seconds);
        $this->timers[$key] = $fiber;
        Fiber::suspend('timer');
    }
}

用起来:

$loop = new Scheduler();

$loop->spawn(function () use ($loop) {
    for ($i = 0; $i < 3; $i++) {
        printf("[%.3f] 慢任务: %dn", microtime(true), $i);
        $loop->sleep(0.3);
    }
});

$loop->spawn(function () use ($loop) {
    for ($i = 0; $i < 10; $i++) {
        printf("[%.3f] 快任务: %dn", microtime(true), $i);
        $loop->sleep(0.05);
    }
});

$loop->run();

慢任务每次睡 300ms,快任务每次睡 50ms,两者交错输出,总耗时约 0.9 秒——而不是各自串行相加的 1.4 秒。

为什么定时器的键要用字符串

这一点很容易踩坑。下面这种写法是错的:

$this->timers[microtime(true) + $seconds] = $fiber;

因为 PHP 数组的浮点键会被隐式转成整数,而微秒部分正好是小数点后的东西,一转就丢。更麻烦的是,转成整数之后,同一秒内创建的多个定时器会撞在同一个键上,前面的直接被后面的覆盖,任务就这么丢了。

用 sprintf('%.6f', ...) 转成字符串键,就能保留六位小数的精度,同时在取最小值的时候用 floatval 转回来比较。

四、第三版:加上流等待

定时器解决的是”等时间”,真正的大头是”等 IO”。这一节加进来的东西,才让这个调度器有了实际价值。

核心思路很简单:一个 Fiber 想读某个 socket 的时候,它把 socket 句柄和”我想读”这个意图登记到调度器,然后挂起自己。调度器在没有就绪任务的时候,用 stream_select 统一等待所有登记的句柄,谁先就绪就把谁对应的 Fiber 唤醒。

final class Scheduler
{
    // ... 前面已有的属性
    /** @var array<int, array{0: Fiber, 1: string}> 资源 id => [Fiber, 方向] */
    private array $streams = [];

    public function waitReadable($stream): void
    {
        $this->park($stream, 'read');
    }

    public function waitWritable($stream): void
    {
        $this->park($stream, 'write');
    }

    private function park($stream, string $mode): void
    {
        $fiber = Fiber::getCurrent();

        // 主线程里直接阻塞等,保证代码在任何上下文都能跑
        if ($fiber === null) {
            $read  = $mode === 'read'  ? [$stream] : [];
            $write = $mode === 'write' ? [$stream] : [];
            $except = null;
            @stream_select($read, $write, $except, null);
            return;
        }

        $this->streams[(int) $stream] = [$fiber, $mode];
        Fiber::suspend('io');
    }

    private function idle(): bool
    {
        if (!$this->timers && !$this->streams) {
            return false;
        }

        $timeout = null;
        if ($this->timers) {
            $next = min(array_map('floatval', array_keys($this->timers)));
            $timeout = max(0.0, $next - microtime(true));
        }

        if ($this->streams) {
            $read = $write = [];
            foreach ($this->streams as $id => [, $mode]) {
                $mode === 'read' ? $read[] = $id : $write[] = $id;
            }
            // stream_select 需要资源,不是 id,这里要映射回去
            $readRes = array_map(fn($id) => $this->resources[$id], $read);
            $writeRes = array_map(fn($id) => $this->resources[$id], $write);
            $except = null;

            if ($timeout === null) {
                $sec = null;
                $usec = null;
            } else {
                $sec = (int) $timeout;
                $usec = (int) (($timeout - $sec) * 1_000_000);
            }

            $n = @stream_select($readRes, $writeRes, $except, $sec, $usec);
            if ($n > 0) {
                foreach ($readRes as $res)  { $this->wake((int) $res); }
                foreach ($writeRes as $res) { $this->wake((int) $res); }
            }
            return true;
        }

        if ($timeout !== null) {
            usleep((int) ($timeout * 1_000_000));
        }
        return true;
    }

    private function wake(int $id): void
    {
        [$fiber] = $this->streams[$id];
        unset($this->streams[$id], $this->resources[$id]);
        $this->ready->enqueue($fiber);
    }
}

上面这段代码有个细节需要处理:stream_select 需要的是资源,而 $streams 用资源 id 做键是为了方便索引。所以在 park 里还要额外存一份 id 到资源的映射:

private array $resources = [];

private function park($stream, string $mode): void
{
    // ...
    $id = (int) $stream;
    $this->streams[$id] = [$fiber, $mode];
    $this->resources[$id] = $stream;
    Fiber::suspend('io');
}

其实还有更简单的做法:直接用资源的 id 去 stream_select?不行,那会报类型错误。PHP 没有提供从资源 id 反查资源的函数,所以这份映射是必要的。

还要注意一点:资源被 fclose 之后,它的 id 可能被系统回收给下一个新建的资源。所以 wake 里必须把两个数组都清干净,否则会留下脏数据,导致下一次唤醒跑错 Fiber。

五、案例:把八个串行请求改成并发

现在用这个调度器写一个非阻塞的 HTTP 抓取器。

final class AsyncHttp
{
    public function __construct(private Scheduler $loop) {}

    public function get(string $url): string
    {
        $parts = parse_url($url);
        $host  = $parts['host'];
        $port  = $parts['port'] ?? 80;
        $path  = ($parts['path'] ?? '/')
               . (isset($parts['query']) ? '?' . $parts['query'] : '');

        $socket = @stream_socket_client(
            "tcp://{$host}:{$port}",
            $errno,
            $errstr,
            5,
            STREAM_CLIENT_CONNECT | STREAM_CLIENT_ASYNC_CONNECT
        );

        if (!$socket) {
            throw new RuntimeException("连接 {$host}:{$port} 失败:{$errstr}");
        }

        stream_set_blocking($socket, false);

        // 异步连接:socket 刚创建时还没连上,等它可写才算连好
        $this->loop->waitWritable($socket);

        $request = "GET {$path} HTTP/1.1rn"
                 . "Host: {$host}rn"
                 . "Connection: closern"
                 . "User-Agent: FiberDemo/1.0rnrn";

        $this->writeAll($socket, $request);

        $raw = $this->readAll($socket);
        fclose($socket);

        // 只做最粗的演示:切掉响应头,保留正文
        $sep = strpos($raw, "rnrn");

        return $sep === false ? $raw : substr($raw, $sep + 4);
    }

    private function writeAll($socket, string $data): void
    {
        $offset = 0;
        $length = strlen($data);

        while ($offset < $length) {
            $written = @fwrite($socket, substr($data, $offset));

            if ($written === false || $written === 0) {
                // 发送缓冲区满了,等可写
                $this->loop->waitWritable($socket);
                continue;
            }

            $offset += $written;
        }
    }

    private function readAll($socket): string
    {
        $buffer = '';

        while (!feof($socket)) {
            $chunk = @fread($socket, 8192);

            if ($chunk === false || $chunk === '') {
                // 当前没有数据可读,让出控制权
                $this->loop->waitReadable($socket);
                continue;
            }

            $buffer .= $chunk;
        }

        return $buffer;
    }
}

发起请求的部分:

$loop = new Scheduler();
$http = new AsyncHttp($loop);

$hosts = ['example.com', 'example.org', 'example.net', 'iana.org'];

$start = microtime(true);

foreach ($hosts as $host) {
    $loop->spawn(function () use ($http, $loop, $host) {
        $t0 = microtime(true);

        try {
            $body = $http->get("http://{$host}/");
            printf("%-16s %7d 字节  耗时 %.3fsn",
                $host, strlen($body), microtime(true) - $t0);
        } catch (Throwable $e) {
            printf("%-16s 失败:%sn", $host, $e->getMessage());
        }
    });
}

$loop->run();

printf("全部完成,总耗时 %.3fsn", microtime(true) - $start);

四个请求会几乎同时发出。stream_socket_client 带 STREAM_CLIENT_ASYNC_CONNECT 的时候立刻返回,连接过程在后台进行,脚本不需要等;随后 waitWritable 把这个 Fiber 挂起,调度器切到下一个任务继续建连接。等到四个连接都建立、或者某个先就绪了,stream_select 会通知调度器唤醒对应的 Fiber 去发请求。

四个响应会按”谁先返回”的顺序陆续打印,而不是按发起顺序。总耗时取决于最慢的那一个,而不是四个之和。

这个 HTTP 客户端能上生产吗

不能。上面这段代码是教学用的,省略了 HTTPS、分块传输编码、重定向、gzip、cookie、超时控制、连接复用等一大堆东西。任何一个真实项目都应该用成熟的客户端——比如 Guzzle 配上支持协程的 handler,或者直接用 Revolt 生态里的 HTTP 客户端。

写它的价值只在于说明一件事:调度器和 IO 之间的接口就这么点东西。想清楚这一点之后,再看 Revolt 或者 AMPHP 的实现,会顺畅很多。

六、七个实际踩过的坑

1. 原生阻塞函数会冻住整个进程

这是最致命的。下面这些函数一旦在 Fiber 里被调用,整个事件循环都会停下来,所有其他任务一起卡住:

sleep()
usleep()
file_get_contents('http://...')
curl_exec()
mysqli_query()  // 同步驱动
PDO::exec()     // 同步驱动

注意 file_get_contents 读本地文件通常很快,可以接受;但读网络 URL 就是灾难。团队里要形成共识:协程代码里出现这些函数必须被代码评审拦下来。

2. 别在析构函数里 suspend

class Connection
{
    public function __destruct()
    {
        Fiber::suspend();  // FiberError: Cannot suspend in a force-closed fiber
    }
}

析构函数的执行时机由 GC 决定,可能发生在任何一个奇怪的位置。在里面挂起会让调度器的状态彻底乱掉。同样的问题也会出现在一个还处于 suspended 状态的 Fiber 对象被垃圾回收的时候——PHP 会报错,而且这个错误很难定位。

实践中的做法是:所有资源清理都写显式的 close() 方法,让任务自己负责在结束前调用。__destruct 只做最后的兜底断言。

3. Fiber 里抛的异常会在调度器里冒出来

// 任务内部
throw new RuntimeException('接口返回 500');

// 结果:run() 里的 $fiber->resume() 直接抛出,
// while 循环被打断,剩下的任务全部不执行

所以 spawn 必须包一层捕获:

public function spawn(callable $task, ?callable $onError = null): void
{
    $this->ready->enqueue(new Fiber(function () use ($task, $onError) {
        try {
            $task();
        } catch (Throwable $e) {
            $onError
                ? $onError($e)
                : fwrite(STDERR, "任务异常:{$e->getMessage()}n");
        }
    }));
}

这样单个任务失败不会拖垮整个循环。如果业务上需要”全部失败就退出”,可以在 $onError 里记录状态,在 run() 结束后统一判断。

4. 全局状态在协程之间是共享的

PHP 的单线程特性意味着,静态变量、$GLOBALS、单例对象在多个 Fiber 之间是同一份。这一点容易产生误解,因为很多人把 Fiber 类比成线程,以为每个 Fiber 有自己的变量空间。

class Auth
{
    public static ?int $userId = null;
}

// 协程 A
Auth::$userId = 101;
$loop->sleep(0.1);
echo Auth::$userId;  // 可能已经被协程 B 改成了 202

解决办法是在切换之前把状态从全局里取出来,存成局部变量,用完了再放回去。或者干脆避免用全局状态,把请求相关的数据都通过参数传递。

5. 数据库连接不能跨协程乱用

PDO 默认是缓冲查询(buffered query),query() 返回的那一刻结果集已经在客户端内存里了,所以多个 Fiber 交替使用同一个连接一般不会出问题。

但如果关掉了缓冲:

$pdo = new PDO($dsn, $user, $pass, [
    PDO::MYSQL_ATTR_USE_BUFFERED_QUERY => false,
]);

那么当一个 Fiber 还在遍历结果集、没读完的时候,另一个 Fiber 发了新查询,MySQL 会直接报 “Cannot execute queries while other unbuffered queries are active”。

还有一种情况更隐蔽:在事务里 suspend。事务开始之后挂起,别的 Fiber 用同一个连接发查询,那些查询就被包进当前事务了。这个 bug 查起来非常痛苦,因为从代码上看两个操作毫无关系。

成熟的方案是给每个 Fiber 分配独立的连接,用完归还连接池。这正是 Swoole、AMPHP 这些框架在底层做的事。

6. getCurrent() 返回的是最内层那个 Fiber

$outer = new Fiber(function () {
    $inner = new Fiber(function () {
        var_dump(Fiber::getCurrent() === $inner);  // true
    });
    $inner->start();
});

$outer->start();

嵌套 Fiber 的时候要有这个意识。虽然这种写法在业务里很少见,但如果你在写框架或者做偏底层的事情,就会遇到。特别是”在调度器里再启动一个临时 Fiber 来做验证”这种场景,很容易不小心让当前 Fiber 的引用变掉。

7. 栈深度和内存没有魔法

每个 Fiber 有自己的调用栈,创建一万个 Fiber 就要一万份栈空间。虽然 PHP 的 Fiber 栈是惰性分配、按需增长的,比线程轻量得多,但也不是免费的。

经验值是几千个并发 Fiber 通常没问题,上万个的时候要留意内存曲线。如果真的有这种量级的需求,更好的做法是用队列把任务分批处理,而不是一次性全开出来。

七、什么时候该用,什么时候别用

写到这里可以给个相对明确的判断了。

适合的场景:脚本里同时存在多种 IO 等待(HTTP + Redis + 文件),而且这些等待的总时长占了绝大部分运行时间。这类场景下 Fiber 的收益非常直接。

不适合的场景,第一个是纯 CPU 密集的任务。任何函数只要在跑计算,就没人能切走,Fiber 完全帮不上忙。想并行算数得用多进程或者扩展。

第二个是”只有一两次 HTTP 请求”的场景。这种直接 curl_multi 更合适——它在 C 层做多路复用,比 PHP 层的事件循环效率高得多,而且不需要引入任何架构上的改动。

第三个是已经在用 Swoole 或者 RoadRunner 的项目。这两个运行时本身就把协程做成了默认行为,你在里面再手写一个调度器只会添乱。

如果确实需要在生产环境用协程,不要用这篇文章里的调度器,去用 Revolt。它是 PHP 官方基金会下的项目,实现了完整的 EventLoop 接口,AMPHP v3 和很多其他库都基于它构建。看它的源码会发现结构跟上面写的很像,只是把各种边界情况都处理干净了——信号处理、fork 安全、多事件循环实例等等。

八、写在最后

回头看这个问题的起点,其实是一个很朴素的观察:程序在等 IO 的时候,CPU 是闲着的。这些年 PHP 生态里出现 Swoole、ReactPHP、AMPHP、Revolt 这一串东西,本质上都在回答同一个问题。

Fiber 的价值在于它把这个能力下沉到了语言层面。以前要写异步代码,得把逻辑拆成回调,或者用 Promise 链,代码的可读性会大打折扣。有了 Fiber 之后,异步代码可以长得跟同步代码一样——顺序写下来,从上往下读,只是中间那些会阻塞的操作换成了调度器提供的方法。

不过也要清楚它的边界。它没让 PHP 变成多线程,没让 CPU 密集任务变快,也没解决共享状态的问题。它解决的只是”等待”。

如果你手上的脚本正好被等待拖慢,那花半天时间把调度器搭起来,收益会很明显。如果不是,那把这份精力放在别的地方更划算。

PHP Fiber 实战:手写一个协程调度器把同步代码跑出并发
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 php PHP Fiber 实战:手写一个协程调度器把同步代码跑出并发 https://www.taomawang.com/server/php/2820.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务