hyperf协程并发处理有五种方法:一、用parallel类控制并发数并聚合结果;二、用channel实现生产-消费模型以支持背压;三、用waitgroup协调无返回值协程完成;四、用协程池复用高频短任务协程;五、用barrier实现多协程同步屏障。

如果您在使用Hyperf处理高并发请求或批量数据时遇到响应延迟、吞吐量瓶颈或协程未按预期并行执行的问题,则很可能是由于协程调度方式不当、阻塞操作未协程化或并发控制策略缺失所致。以下是实现Hyperf协程并发处理的多种方法:
一、使用 Parallel 工具类进行可控并发
Parallel 是 Hyperf 提供的高层并发封装,适用于需限制并发数、等待全部结果返回的场景,如批量接口调用、分片数据处理等。它内部基于协程创建与 WaitGroup 机制,确保资源可控且不压垮服务。
1、在控制器或服务中引入 Parallel 类:use Hyperf\Utils\Parallel;
2、实例化 Parallel 并设定最大并发数(例如 50):$parallel = new Parallel(50);
3、通过 add 方法注册每个子任务闭包,闭包内可执行协程安全的 I/O 操作:$parallel->add(function() use ($id) { return $this->fetchUserData($id); });
4、调用 wait() 阻塞至所有协程完成,并获取聚合结果数组:$results = $parallel->wait();
二、手动创建协程配合 Channel 实现生产-消费模型
当需要解耦任务生成与执行节奏、支持背压控制或异步缓冲时,可使用 Swoole\Coroutine\Channel 构建协程通道,实现多生产者向单/多消费者投递任务的模式,避免内存溢出和协程失控。
1、初始化一个有界通道,容量设为合理上限(如 100):$chan = new Swoole\Coroutine\Channel(100);
2、启动多个生产协程,向通道推送任务数据:Swoole\Coroutine::create(function() use ($chan) { for ($i = 0; $i push(['uid' => $i, 'type' => 'profile']); } });
3、启动固定数量消费者协程(如 5 个),持续从通道拉取并处理:for ($j = 0; $j pop()) && ! $chan->isClosed()) { $this->processUserTask($task); } }); }
4、所有生产协程结束后关闭通道:$chan->close();
三、利用 WaitGroup 协调无序协程组的完成状态
WaitGroup 适用于无需返回值、仅需确保若干协程全部结束再继续后续逻辑的场景,例如日志上报、缓存预热、异步通知等。它不依赖通道,轻量且语义清晰。
1、实例化 WaitGroup 对象:$wg = new Swoole\Coroutine\WaitGroup();
2、在每个协程启动前调用 add() 增加计数:$wg->add();
3、在协程内部执行完任务后调用 done() 减少计数:Swoole\Coroutine::create(function() use ($wg) { $this->sendNotification(); $wg->done(); });
4、主协程调用 wait() 阻塞直至所有子协程完成:$wg->wait();
四、基于协程池(Coroutine Pool)管理高频短任务
对于频繁创建销毁协程开销敏感的场景(如每秒数千次数据库查询),应复用协程而非每次 new,协程池通过预分配与回收机制降低调度成本,提升稳定性。
1、定义协程池工厂类,实现 make() 和 close() 方法:class TaskPoolFactory implements PoolFactoryInterface { public function make(): object { return new TaskExecutor(); } }
2、在 config/autoload/pool.php 中注册池配置,设置最小/最大协程数:'min_connections' => 10, 'max_connections' => 200
3、通过容器获取协程池实例:$pool = $this->container->get(CoroutinePool::class);
4、使用 get() 获取可用协程对象执行任务,完成后归还:$executor = $pool->get(); $executor->run($data); $pool->put($executor);
五、结合 Swoole\Coroutine\Barrier 实现多协程同步屏障
Barrier 用于让一组协程在指定检查点集体等待,直到全部到达后再统一继续执行,适用于需要严格时序对齐的并发流程,如分布式锁预检、多源数据比对起始同步点等。
1、创建 Barrier 实例并传入预期协程总数(如 8):$barrier = new Swoole\Coroutine\Barrier(8);
2、每个协程执行到同步点时调用 await() 进入等待:Swoole\Coroutine::create(function() use ($barrier) { $this->prepareData(); $barrier->await(); });
3、当第 8 个协程调用 await() 后,所有协程同时被唤醒继续执行后续逻辑。











