今天是第一周周三,我们将深入 Swoole 协程最核心的工具箱:Channel(通道)、WaitGroup(等待组)、协程上下文(Context)和 defer(延迟执行)。这些是编写健壮、高性能协程应用的基石,也是 Hyperf 框架内部大量使用的机制。学完今天,你将能够自信地构建复杂的协程工作流,并写出第一个协程安全的单元测试。


在这里插入图片描述

今日目标

  1. 掌握 Swoole\Coroutine\Channel 的使用,实现线程间安全的生产者-消费者模式。
  2. 使用 Swoole\Coroutine\WaitGroup 协调多个并发协程,等待其全部完成。
  3. 理解协程上下文 Context 的作用域,并能在 HTTP 请求中安全传递请求 ID。
  4. 学会使用 defer 进行资源清理,确保文件句柄、锁等被正确释放。
  5. 编写 PHPUnit 协程测试,验证并发环境下的数据一致性。

一、环境准备与工具安装 (约 30 分钟)

进入我们的一贯工作环境,并准备好 PHPUnit 以便进行单元测试。

1. 进入容器并创建今日目录
cd swoole-course
docker-compose exec swoole bash
mkdir -p day3
cd day3
2. 安装 PHPUnit(在容器内)

由于 Swoole 官方镜像不包含 PHPUnit,我们手动安装它。

# 下载 phpunit.phar 到当前目录
wget -O phpunit https://phar.phpunit.de/phpunit-10.phar
chmod +x phpunit
# 移动到可执行路径,或者直接使用相对路径
# 建议直接在当前目录使用 ./phpunit

验证安装:

./phpunit --version

现在我们就有了测试利器。


二、知识核心:协程工具箱深度解析 (约 1.5 小时)

1. Channel(通道)—— 协程间的安全桥梁

类比:Channel 像一个线程安全的管道或队列,生产者把数据 push 进去,消费者 pop 出来。它完美解决了共享内存的加锁问题。

特性

  • 容量 (capacity):若指定容量,通道满时 push 会挂起当前协程,直到有空间;通道空时 pop 也会挂起,直到有数据。天然的流量控制。
  • 无需锁:底层基于原子操作和协程调度,保证线程安全。
  • 可用于通信:不仅仅传递数据,还可传递信号(如关闭信号)。

常用方法

$chan = new Swoole\Coroutine\Channel(10); // 容量为 10 的通道
$chan->push($data);   // 写入数据
$data = $chan->pop(); // 读取数据(若空则协程挂起)
$chan->close();       // 关闭通道,所有 pop 会立刻返回 false
2. WaitGroup —— 优雅地等待一群协程

类比:就像餐厅里,服务员需要等厨房里所有菜都做好才一起上桌。WaitGroup 就是那个计数器。

基本用法

$wg = new Swoole\Coroutine\WaitGroup();
$wg->add(); // 计数器 +1
go(function () use ($wg) {
    // 执行任务
    $wg->done(); // 任务完成,计数器 -1
});
$wg->wait(); // 等待计数器归零
3. 协程上下文 (Context) —— 协程专属的全局变量

在每个协程中,我们可能需要传递一些数据,如请求 ID、用户信息等,但又不想弄脏函数参数。Context 提供了一个隔离的键值存储,每个协程拥有自己独立的上下文,子协程可以继承父协程的上下文。

  • Swoole\Coroutine\Context::get($key, $default = null)
  • Swoole\Coroutine\Context::set($key, $value)

在 Hyperf 中,我们常用 Context::set/get,但其底层原理相同。

4. defer —— 延迟执行,资源清理利器

Swoole\Coroutine::defer(callable $callback):注册一个回调,在当前协程结束前一定会执行。即使协程异常退出,也会执行。非常适合释放锁、关闭文件、归还连接池等。


三、实战:构建协程工作流 (约 2.5 小时)

实战 1:生产者-消费者爬虫模拟(Channel 实战)

我们设计一个场景:多个协程生产 URL,另几个协程消费这些 URL 并“爬取”(模拟耗时 0.5 秒),最后汇总结果。这完美展示了 Channel 的流量控制与 WaitGroup 的等待。

创建文件 producer_consumer.php

<?php
use Swoole\Coroutine;
use Swoole\Coroutine\Channel;
use Swoole\Coroutine\WaitGroup;

// 启动时设置一个全局协程化环境
Co\run(function () {
    $start = microtime(true);
    $urls = [
        'http://site1.com', 'http://site2.com', 'http://site3.com',
        'http://site4.com', 'http://site5.com', 'http://site6.com',
    ];
    
    // 1. 创建通道,容量为 2,故意设小一点观察流量控制
    $chan = new Channel(2);
    
    // 2. 用于等待消费者完成
    $wg = new WaitGroup();
    
    // 3. 启动生产者协程
    Coroutine::create(function () use ($chan, $urls) {
        foreach ($urls as $url) {
            echo "[生产者] 推送: {$url}\n";
            $chan->push($url);
        }
        // 所有 URL 生产完毕,关闭通道以通知消费者退出
        $chan->close();
        echo "[生产者] 生产完毕,通道已关闭\n";
    });
    
    // 4. 启动多个消费者协程 (3个)
    $results = []; // 注意:多个协程写同一个数组不安全,仅演示用
    for ($i = 1; $i <= 3; $i++) {
        $wg->add(); // 计数器+1
        $consumerId = $i;
        Coroutine::create(function () use ($chan, $wg, $consumerId, &$results) {
            while (true) {
                $url = $chan->pop(); // 如果通道空且未关闭,会挂起
                if ($url === false) {
                    // 通道已关闭且无数据
                    break;
                }
                echo "[消费者{$consumerId}] 正在爬取: {$url}\n";
                // 模拟爬取耗时
                Coroutine\System::sleep(0.5);
                $results[] = "{$url} - 200 OK";
                echo "[消费者{$consumerId}] 完成: {$url}\n";
            }
            $wg->done(); // 任务完成
            echo "[消费者{$consumerId}] 退出\n";
        });
    }
    
    // 5. 等待所有消费者完成
    $wg->wait();
    
    // 打印结果
    echo "所有爬取结果:\n";
    var_dump($results);
    echo "总耗时: " . (microtime(true) - $start) . " 秒\n";
    // 注意:并发执行,总耗时约 0.5秒 * ceil(6/3) = 1秒,而非 3秒
});

运行观察

php producer_consumer.php

你会看到生产者推送时,当通道满(2个)时,push 会挂起,直到有消费者取走数据。这自动实现了背压。

实战 2:WaitGroup 与请求上下文传递(与 HTTP 结合)

我们结合周二的 HTTP 服务器,实现每个请求生成唯一的 request_id,并通过 Context 在整个协程生命周期内传递,日志中打印出来。

创建文件 http_context.php

<?php
use Swoole\Http\Server;
use Swoole\Http\Request;
use Swoole\Http\Response;
use Swoole\Coroutine;
use Swoole\Coroutine\Context;

$server = new Server('0.0.0.0', 9501);

$server->on('Request', function (Request $req, Response $res) {
    // 1. 为当前请求生成唯一 ID,并存入协程上下文
    $requestId = uniqid('req_', true);
    Context::set('request_id', $requestId);
    
    // 2. 记录请求开始日志
    logMessage("请求开始: {$req->server['request_uri']}");
    
    // 3. 模拟一个需要调用多个下游服务的复杂操作,使用 WaitGroup 协调
    $wg = new Coroutine\WaitGroup();
    $userData = $orderData = null;
    
    // 子协程1:获取用户信息(耗时0.2秒)
    $wg->add();
    Coroutine::create(function () use ($wg, &$userData) {
        // 子协程会继承父协程的上下文,可以拿到 request_id
        logMessage("正在获取用户信息...");
        Coroutine\System::sleep(0.2);
        $userData = ['id' => 1, 'name' => 'Swoole'];
        logMessage("用户信息获取完成");
        $wg->done();
    });
    
    // 子协程2:获取订单信息(耗时0.3秒)
    $wg->add();
    Coroutine::create(function () use ($wg, &$orderData) {
        logMessage("正在获取订单信息...");
        Coroutine\System::sleep(0.3);
        $orderData = ['order_id' => 1001, 'amount' => 99.9];
        logMessage("订单信息获取完成");
        $wg->done();
    });
    
    // 4. 等待所有子协程完成
    $wg->wait();
    
    // 5. 组装响应
    $result = [
        'code' => 200,
        'request_id' => $requestId,
        'data' => [
            'user' => $userData,
            'order' => $orderData,
        ]
    ];
    
    logMessage("请求处理完成");
    $res->header('Content-Type', 'application/json');
    $res->end(json_encode($result));
});

function logMessage(string $msg) {
    $requestId = Context::get('request_id', 'unknown');
    echo "[{$requestId}] " . $msg . PHP_EOL;
}

$server->start();

启动并测试

php http_context.php

使用 curl http://localhost:9501/ 发送几个请求,你会看到控制台输出中,每个请求的日志都带有自己的 request_id,互不干扰。这就是协程上下文隔离的强大之处。

实战 3:defer 清理资源与锁释放

假设我们有一个模拟的“文件锁”或者“连接资源”,需要在使用后安全释放。我们用 defer 来确保清理。

创建文件 defer_demo.php

<?php
use Swoole\Coroutine;

Coroutine\run(function () {
    $file = 'demo.txt';
    
    // 协程1:写入文件,并使用 defer 保证关闭
    go(function () use ($file) {
        $fp = fopen($file, 'w');
        if (!$fp) {
            echo "无法打开文件\n";
            return;
        }
        
        // 注册 defer,无论协程如何结束,都会执行此回调
        Coroutine::defer(function () use ($fp) {
            echo "关闭文件句柄\n";
            fclose($fp);
        });
        
        fwrite($fp, "Hello Swoole\n");
        Coroutine\System::sleep(1);
        fwrite($fp, "Another line\n");
        
        // 即使我们抛出异常,defer 也会执行
        // throw new \Exception('some error');
        
        echo "正常写入完成\n";
    });
    
    // 协程2:模拟获取分布式锁,用 defer 释放
    go(function () {
        $lockKey = 'resource_lock';
        echo "尝试获取锁...\n";
        // 模拟获取锁成功
        $locked = true;
        
        Coroutine::defer(function () use ($lockKey) {
            echo "释放锁: {$lockKey}\n";
        });
        
        if ($locked) {
            echo "执行业务逻辑...\n";
            Coroutine\System::sleep(0.5);
            echo "业务完成\n";
        }
        
        // 协程结束,defer 自动执行
    });
});

运行 php defer_demo.php,你会发现无论中间发生什么,defer 中的释放代码都会在最后被调用。注意:defer 是绑定在当前协程上的,不同协程的 defer 互相独立。


四、成果测试:编写协程单元测试 (约 1.5 小时)

PHPUnit 结合 Swoole 协程环境,我们需要使用 Swoole\Coroutine\run 包裹测试用例,或者直接使用 Hyperf 测试框架(但那在后面)。今天我们先裸写一个测试来验证生产者-消费者数据一致性。

1. 配置 PHPUnit

新建 phpunit.xmlday3 目录:

<?xml version="1.0" encoding="UTF-8"?>
<phpunit bootstrap="vendor/autoload.php" colors="true">
    <testsuites>
        <testsuite name="SwooleCoroutine">
            <directory>tests</directory>
        </testsuite>
    </testsuites>
</phpunit>

创建 tests 目录:

mkdir tests
2. 编写生产者-消费者测试

创建 tests/ChannelTest.php

<?php
use PHPUnit\Framework\TestCase;
use Swoole\Coroutine;
use Swoole\Coroutine\Channel;
use Swoole\Coroutine\WaitGroup;

class ChannelTest extends TestCase
{
    public function testProducerConsumerDataConsistency()
    {
        // 在协程环境中运行测试
        Coroutine\run(function () {
            $totalItems = 100;
            $chan = new Channel(10);
            $consumed = 0;
            $wg = new WaitGroup();
            
            // 生产者
            Coroutine::create(function () use ($chan, $totalItems) {
                for ($i = 1; $i <= $totalItems; $i++) {
                    $chan->push($i);
                }
                $chan->close();
            });
            
            // 消费者(2个)
            for ($c = 0; $c < 2; $c++) {
                $wg->add();
                Coroutine::create(function () use ($chan, $wg, &$consumed) {
                    while (true) {
                        $item = $chan->pop();
                        if ($item === false) {
                            break;
                        }
                        // 模拟处理
                        $consumed++;
                    }
                    $wg->done();
                });
            }
            
            $wg->wait();
            
            // 断言:消费总数应等于生产总数
            $this->assertEquals($totalItems, $consumed, "数据丢失或重复!");
            echo "测试通过: 生产 {$totalItems},消费 {$consumed}\n";
        });
    }
    
    public function testMultipleProducersSingleConsumer()
    {
        Coroutine\run(function () {
            $chan = new Channel(5);
            $total = 0;
            $wg = new WaitGroup();
            
            // 3个生产者,各生产 50 个
            for ($p = 0; $p < 3; $p++) {
                $wg->add();
                Coroutine::create(function () use ($chan, $wg) {
                    for ($i = 0; $i < 50; $i++) {
                        $chan->push(1);
                    }
                    $wg->done();
                });
            }
            
            // 消费者
            $consumerWg = new WaitGroup();
            $consumerWg->add();
            Coroutine::create(function () use ($chan, $consumerWg, &$total) {
                while (true) {
                    $val = $chan->pop();
                    if ($val === false) break;
                    $total++;
                }
                $consumerWg->done();
            });
            
            // 等待所有生产者完成,然后关闭通道
            $wg->wait();
            $chan->close();
            $consumerWg->wait();
            
            $this->assertEquals(150, $total);
        });
    }
}
3. 运行测试
./phpunit tests/ChannelTest.php

如果一切正常,你应该看到绿色的通过提示。这就证明了我们的协程操作是线程安全的。

进阶:可以故意修改代码(例如不加锁并发写数组),看看测试是否会失败,深刻理解 Channel 的重要性。


五、今日检验清单

知识点检验方法通过标准
Channel 通信运行 producer_consumer.php所有 URL 被消费,通道满时生产者阻塞,结果无丢失
WaitGroup 协调HTTP 服务器返回包含用户和订单的完整数据请求总耗时约 0.3s(并发),而非 0.5s
Context 隔离同时发送多个 HTTP 请求,观察日志每个请求日志带有唯一 request_id,不窜乱
defer 资源释放运行 defer_demo.php,即使抛出异常也能看到释放日志“关闭文件句柄”、“释放锁” 必定出现在最后
协程单元测试执行 PHPUnit 测试全部绿灯,数据一致性断言通过

今日作业与学习产出

  1. 提交代码:将 producer_consumer.phphttp_context.phpdefer_demo.php 以及 tests/ 目录提交到 Git。
  2. 学习笔记:画一个生产者-消费者的协程时序图,标出 Channel push/pop 的阻塞唤醒时机。
  3. 挑战任务:修改生产者-消费者代码,实现动态退出:当生产者生产完所有 URL 后,消费者还能从 Channel 接收到一个特殊的结束信号(如字符串 "EOF"),再退出循环。对比 close() 的方式,哪种更好?
  4. 挑战任务:在 HTTP 上下文传递的基础上,添加 defer 来记录每个请求的结束时间(在 defer 中计算耗时并输出)。

通过今天的课程,你已经掌握了构建复杂协程应用的必备工具,为明天正式迈入 Hyperf 框架打下了最坚实的基础。

更多推荐