如何在Workerman循环内创建多进程?实现数组并行遍历
在Workerman中实现多进程遍历数组的方案
首先明确:你原代码的写法不可行,问题出在这两点:
$worker->count是Worker启动前的配置参数,启动后在onWorkerStart里修改完全无效,无法动态创建新进程。- 循环里的逻辑还是单进程串行执行,根本没用到多进程并行能力。
下面提供三种可行的实现方案,按需选择:
方案一:基于Worker多进程的分片处理(最简单)
直接启动指定数量的Worker进程,将数组拆分为多个分片,每个进程负责处理自己的分片,适合任务耗时均匀的场景。
use Workerman\Worker; require_once __DIR__ . '/vendor/autoload.php'; // 待处理的数组 $targetArray = range(1, 100); // 设定进程数量 $processNum = 10; // 拆分数组为对应进程数的分片 $arrayChunks = array_chunk($targetArray, ceil(count($targetArray)/$processNum)); $worker = new Worker(); // 启动指定数量的进程 $worker->count = $processNum; $worker->onWorkerStart = function($worker) use ($arrayChunks) { // 每个进程根据自身ID获取对应的数组分片 $currentChunk = $arrayChunks[$worker->id]; foreach($currentChunk as $item) { // 替换为你的业务处理逻辑 echo "进程{$worker->id}处理元素: {$item}\n"; // 模拟业务耗时 sleep(1); } // 当前进程任务完成后退出 Worker::stopAll(); }; Worker::runAll();
方案二:手动创建Process进程(更灵活)
通过Workerman的Process类手动创建子进程,可按需控制进程的创建时机和数量,适合动态调整进程数的场景。
use Workerman\Worker; use Workerman\Process; require_once __DIR__ . '/vendor/autoload.php'; $targetArray = range(1, 100); $processNum = 10; $arrayChunks = array_chunk($targetArray, ceil(count($targetArray)/$processNum)); // 主进程仅负责创建子进程 $mainWorker = new Worker(); $mainWorker->onWorkerStart = function() use ($arrayChunks) { foreach($arrayChunks as $index => $chunk) { // 为每个数组分片创建独立进程 $process = new Process(function() use ($index, $chunk) { foreach($chunk as $item) { echo "进程{$index}处理元素: {$item}\n"; // 替换为你的业务逻辑 sleep(1); } }); $process->start(); } // 等待所有子进程执行完毕 Process::wait(); // 主进程退出 Worker::stopAll(); }; Worker::runAll();
方案三:基于Channel组件的动态任务分发(负载均衡)
当数组元素处理耗时差异较大时,用Channel组件实现主进程分发任务、子进程抢任务的模式,避免进程负载不均。
use Workerman\Worker; use Workerman\Channel\Channel; require_once __DIR__ . '/vendor/autoload.php'; $targetArray = range(1, 100); $processNum = 10; $totalTasks = count($targetArray); $completedTasks = 0; // 启动Channel服务,用于进程间通信 $channel = new Channel(); $channel->listen(); // 任务分发进程 $dispatchWorker = new Worker(); $dispatchWorker->onWorkerStart = function() use ($targetArray, $processNum) { // 等待所有处理进程启动 sleep(1); // 分发所有任务 foreach($targetArray as $item) { Channel::publish('task_queue', $item); } // 给每个处理进程发送结束信号 for($i=0; $i<$processNum; $i++) { Channel::publish('task_queue', null); } }; // 任务处理进程 $worker = new Worker(); $worker->count = $processNum; $worker->onWorkerStart = function($worker) use (&$completedTasks, $totalTasks) { // 订阅任务频道 Channel::subscribe('task_queue', function($item) use ($worker, &$completedTasks, $totalTasks) { if(is_null($item)) { // 收到结束信号,进程退出 $worker->exit(); return; } // 业务处理逻辑 echo "进程{$worker->id}处理元素: {$item}\n"; sleep(1); // 统计完成任务数 $completedTasks++; if($completedTasks >= $totalTasks) { // 所有任务完成,停止所有进程 Worker::stopAll(); } }); }; Worker::runAll();
总结:你的需求完全可以实现,只要避开原代码的错误写法,根据任务特点选择上述方案即可。
内容的提问来源于stack exchange,提问作者Евгений Антипов
相关产品推荐
相关产品推荐

