PHP中如何进行多任务协作调度
一、什么是多任务协作
二、进程,线程,协程
三、PHP通过多进程的方式来实现多任务协作
四、PHP通过协程的方式来实现多任务协作
五、结语
一、什么是多任务协作
简单来说,就是在一个程序内,多个任务协同执行,并可以相互联系也可以没有任何的依赖。举一个简单的例子,在业务需求中,我们需要执行一个跑数任务,这个跑数任务可以拆分为三个子任务,主要功能是:
1、在Hive中执行一些时间较长的计算,拿到三个月的指定数据(假如由于业务特殊性需要跑三次SQL)
2、在CK中执行一些相关的跑数计算
3、将前两个任务的数据拿到之后进行处理,再导入到指定的Holo表中
我们可以看到,在任务1中,需要跑三次才能拿到数据,并且我们都知道在Hive中执行一个大的SQL语句耗时是非常久的,更别说需要跑三次,那么在这种情况下,为了节省跑数时间,我们希望这些子任务能够同时进行,按照指定的顺序去组合得到最终的结果,提高程序的执行效率,这就是多任务之间的协作与调度。
当然这个例子举的很苍白,在实际业务中使用也需要小心翼翼,因为并行执行几个大型的SQL可能不是一个好的解决方案。
多任务协作的实现方式有太多太多,例如Golang中的协程goroutine,python中的yield,java中的thread,甚至于nodejs的异步编程等,但我们今天讨论的主要是如何使用PHP去实现多任务协作。
二、进程,线程,协程
在了解如何实现之前,我们需要清楚几个概念,什么是进程,线程和协程,他们之间有什么区别?
1、进程,相信大家都很熟悉了
进程是程序一次动态执行的过程,是程序运行、系统资源分配的基本单位
每个进程都有自己的独立内存空间,不同进程通过进程间通信来通信
进程占据独立的内存,所以上下文进程间的切换开销(栈、寄存器、页表、文件句柄等)比较大,但相对比较稳定安全
2、线程
线程又叫做轻量级进程,是CPU调度的最小单位
线程从属于进程,是程序的实际执行者。一个进程至少包含一个主线程,也可以有更多的子线程
多个线程共享所属进程的资源,同时线程也拥有自己的专属资源
程间通信主要通过共享内存,上下文切换很快,资源开销较少,但相比进程不够稳定容易丢失数据
3、协程
协程是一种用户态的轻量级线程,协程的调度完全由用户控制
一个线程可以拥有多个协程,协程不是被操作系统内核所管理,而完全是由程序所控制
与其让操作系统调度,不如让自己来,这就是协程,goroutine以及swoole的协程都是如此
创建与销毁的带来的开销极小
4、线程与进程的区别:
地址空间:线程是进程内的一个执行单元,进程内至少有一个线程,它们共享进程的地址空间,而进程有自己独立的地址空间
资源拥有:进程是资源分配和拥有的单位,同一个进程内的线程共享进程的资源
线程是处理器调度的基本单位,但进程不是
每个独立的线程有一个程序运行的入口、顺序执行序列和程序的出口,但是线程不能够独立执行,必须依存在应用程序中,由应用程序提供多个线程执行控制
5、协程与线程的区别:
一个线程可以多个协程,一个进程也可以单独拥有多个协程
线程进程都是同步机制,而协程则是异步
协程能保留上一次调用时的状态,每次过程重入时,就相当于进入上一次调用的状态
线程是抢占式,而协程是非抢占式的,所以需要用户自己释放使用权来切换到其他协程,因此同一时间其实只有一个协程拥有运行权,相当于单线程的能力
三、PHP通过多进程的方式来实现多任务协作
了解进程、线程、协程的概念之后,我们开始使用PHP以多进程的方式来实现多任务协作。使用进程的方式去实现多任务是最简单的,但同时也是最容易出问题的,为什么会这么说呢,我们一步一步来实现。
首先我们定义一个 Scheduler 类,用于对任务进行添加以及执行等操作
Scheduler
class Scheduler{protected $maxTaskId = 0;//任务映射protected $taskMap = []; // taskId => [task=>$task,pid=>$pid]//任务队列protected $taskQueue;//主进程的pidprotected $pid;//任务结果集protected $taskResult = [];public function __construct(){$this->taskQueue = new SplQueue();$this->pid = posix_getpid();}/*** 添加一个任务* @param callable $task* @return int*/public function addTask(callable $task){$tid = ++$this->maxTaskId;$this->taskMap[$tid]['task'] = $task;$this->schedule($task, $tid);return $tid;}/*** 调度到队列里* @param $task* @param $id* @return void*/public function schedule($task, $id){$this->taskQueue->enqueue(['id' => $id, 'task' => $task]);}/*** 开始执行任务* @return void*/public function run(){try {$execute_time = time();while (!$this->taskQueue->isEmpty()) {//任务出队$task = $this->taskQueue->dequeue();//fork一个子进程$pid = pcntl_fork();//记录子进程的PID$this->taskMap[$task['id']]['pid'] = $pid;if ($pid == -1) {echo '创建子进程失败' . PHP_EOL;exit(0);} elseif ($pid == 0) {//子进程$result = call_user_func($task['task']);$msg = ['data' => $result,'task_id' => $task['id']];exit();}}//todo 执行等待任务结束后返回结果的操作} catch (Exception $exception) {echo '[执行任务错误]' . $exception->getMessage() . PHP_EOL;exit();}}public function getResult($task_id = 0){if ($task_id != 0) {return $this->taskResult[$task_id];}return $this->taskResult;}}
通过Scheduler类,我们可以通过addTask()的方法来添加新的任务,并且通过run()方法去执行这些任务,原理也很简单,用了SplQueued队列来实现任务的投递,在run时取出任务并且创建一个子进程。
定义完调度类之后,我们需要先生成一些demo代码,来模拟多任务协作的方式。下面代码用sleep来模拟不同任务完成的时间:
demo
$schedule = new Scheduler();//mock任务function mockTask($sleep){sleep($sleep);echo "任务执行了{$sleep}秒" . PHP_EOL;}$task1 = function () {mockTask(10);return 'task1';};$task2 = function () {mockTask(4);return 'task2';};//添加任务$schedule->addTask($task1);$schedule->addTask($task2);$schedule->run();var_dump($schedule->getResult());exit();
但是这样的Scheduler类还存在着一些问题,首先,就是创建完子进程之后,主进程就退出了,没有任何的阻塞,这样我们对子进程则是一无所知,只能通过 ps -ef |grep xxxx 来查看子进程是否执行完了,第二个问题,子进程虽然执行完了,但是由于父子进程之间缺少通信,所以我们并不能知道子进程的执行结果。思路就是,通过一个死循环的方式,等待所有的子进程执行结束,再退出。那么子进程如何通知主进程它执行完了呢。其实有很多种方式实现,例如消息队列(Redis、Kafka),双通道socket的流处理(可以参考使用stream_socket_pair函数来实现),以及信号量来处理。我们这里使用的是共享内存片段来进行通讯,通过创建一个共享内存片段,对子进程的执行结果进行返回。
针对上面所说的问题,我们可以对代码做一些修改
public function __construct(){$this->taskQueue = new SplQueue();$this->pid = posix_getpid();//类初始化时就要设置一个共享内存块了$this->shareMemorySetting();}/*** 共享内存块设置* @return void*/public function shareMemorySetting(){if (function_exists("shm_attach") === FALSE) {die("\n环境不支持shm_attach,前往http://us2.php.net/manual/en/shmop.setup.php安装");}//生成shmKey(实际使用中要考虑有可能产生key碰撞问题)$shmKey = intval(rand(1, 9) . time() . rand(1, 9999));//共享内存大小,可以设置为用户可控制$size = 1024 * 1024 * 10;//开出一片内存空间$this->shareMemory = shm_attach($shmKey, $size);}/*** 开始执行任务* @return void*/public function run(){try {$execute_time = time();while (!$this->taskQueue->isEmpty()) {$task = $this->taskQueue->dequeue();$pid = pcntl_fork();$this->taskMap[$task['id']]['pid'] = $pid;if ($pid == -1) {echo '创建子进程失败' . PHP_EOL;exit(0);} elseif ($pid == 0) {//子进程$result = call_user_func($task['task']);$msg = ['data' => $result,'task_id' => $task['id']];shm_put_var($this->shareMemory, intval($task['id']), $msg);exit();}}while (true) {foreach ($this->taskMap as $task_id => $task_list) {if (shm_has_var($this->shareMemory, $task_id) && !isset($this->taskResult[$task_id])) {$this->taskResult[$task_id] = shm_get_var($this->shareMemory, $task_id)['data'];}}if (count($this->taskMap) <= count($this->taskResult)) {break;}}} catch (Exception $exception) {echo '[执行任务错误]' . $exception->getMessage() . PHP_EOL;exit();}}//类被销毁的时候一定要把共享内存块删除public function __destruct(){shm_remove($this->shareMemory);}
第一个while循环用于创建子进程,并且子进程执行任务,把处理结果塞到共享内存。第二个while循环用于等待处理结果,直到最后一个任务执行完毕才退出。当然,我们还需要加上任务执行超时的操作,以及kill任务等控制任务的操作,避免任务执行时间太长或其他原因导致可能内存溢出等风险。
protected $maxWaitSec = 60;/*** 设置等待时间 0 为无限等待* @return $this*/public function setMaxWaitSec($sec){$this->maxWaitSec = $sec;return $this;}/*** 开始执行任务* @return void*/public function run(){try {$execute_time = time();while (!$this->taskQueue->isEmpty()) {$task = $this->taskQueue->dequeue();$pid = pcntl_fork();$this->taskMap[$task['id']]['pid'] = $pid;if ($pid == -1) {echo '创建子进程失败' . PHP_EOL;exit(0);} elseif ($pid == 0) {//子进程$result = call_user_func($task['task']);$msg = ['data' => $result,'task_id' => $task['id']];shm_put_var($this->shareMemory, intval($task['id']), $msg);exit();}}while (true) {//执行超时的判断 0为无限制等待if ($this->maxWaitSec != 0 && time() - $execute_time > $this->maxWaitSec) {$this->killTask();throw new Exception('执行超时');}foreach ($this->taskMap as $task_id => $task_list) {if (shm_has_var($this->shareMemory, $task_id) && !isset($this->taskResult[$task_id])) {$this->taskResult[$task_id] = shm_get_var($this->shareMemory, $task_id)['data'];}}if (count($this->taskMap) <= count($this->taskResult)) {break;}}} catch (Exception $exception) {echo '[执行任务错误]' . $exception->getMessage() . PHP_EOL;exit();}}public function killTask($tid = 0){$kill_list = [];if ($tid == 0) {$kill_list = array_column($this->taskMap, 'pid');$this->taskMap = [];} else {if (!isset($this->taskMap[$tid])) {return false;}$kill_list[] = $this->taskMap[$tid]['pid'];unset($this->taskMap[$tid]);}foreach ($kill_list as $pid) {exec("kill -9 $pid");}return true;}
执行任务的demo代码如下:
$schedule = new Scheduler();function mockTask($sleep){sleep($sleep);echo "任务执行了{$sleep}秒" . PHP_EOL;}$task1 = function () {mockTask(10);return 'task1';};$task2 = function () {mockTask(4);return 'task2';};//以下涉及到的db操作类来源于天翼$db = DbFactory::getInstance('holo');$db->getPdo()->setAttribute(PDO::ATTR_EMULATE_PREPARES, true);$sqlTask1 = function () use ($db) {$sql = "SELECT * FROM bi_middle_platform.platform_sub_game_info limit 1";return $db->query($sql)->fetchAll(PDO::FETCH_ASSOC);};//设置超时等待时间$schedule->setMaxWaitSec(30);$schedule->addTask($task1);$schedule->addTask($task2);$schedule->addTask($sqlTask1);$schedule->run();foreach ($schedule->getResult() as $task_id => $result) {echo "任务{$task_id}执行结果:" . json_encode($result, JSON_UNESCAPED_UNICODE) . PHP_EOL;}exit();
执行的结果如下:
根据打印的顺序,我们可以知道,最开始执行完成的是sqlTask1这个任务,然后是task2和task1,并且我们分别拿到了他们对应的结果,等待的时间也应该是最长的执行时间,也就是10秒。
以上就是PHP利用多进程的方式来简单实现多任务的协作。在这里,为了演示需要,通常以父进程为主,父进程退出,子进程也会跟着退出,实际上,即使在父进程结束后,子进程也可以继续运行。或者孩子可以杀死父母。可以修改调度程序以具有更分层的任务结构,
还有更多可以实现的流程管理调用以及安全措施,例如并发锁机制,父子进程信号量的监听等。
四、PHP通过协程的方式来实现多任务协作
1、安装swoole扩展来实现协程
如何通过PHP来实现像Golang那样的协程呢?已经有大神为我们开发出了swoole这一协程框架。swoole是一款高性能的PHP协程框架,以扩展的形式使用,安装它之后,便可以使用Swoole空间下的相关协程函数,当然swoole不仅仅有协程,也有websocket、Http、Redis服务器、毫秒级定时器、事件管理等强大的功能。那么如何安装呢?如果你的电脑上有pecl的话,可以直接通过 pecl install swoole 进行安装,也可以参考官方的文档:
swoole官方文档 https://www.swoole.com/
相信了解过swoole这个扩展的同学都知道,swoole与golang一样,通过"go"关键字开启一个协程,与golang一样,在swoole扩展的加持下,在PHP中实现协程变得简单且触手可及,以下是swoole创建100个协程的一段示例代码:
swoole实现协程
Co::set(['hook_flags' => SWOOLE_HOOK_TCP]);Co\run(function() {for ($c = 10; $c--;) {go(function () {//创建10个协程$sleep = rand(1, 10);Coroutine::sleep($sleep); //此处产生协程调度,cpu切到下一个协程,不会阻塞进程,注意不能直接用sleep,直接用sleep会阻塞,要用官方给的Sleep函数echo "我是第{$c}个创建出来的协程,我sleep了「{$sleep}」秒" . PHP_EOL;});}});
运行两次结果如下:
我们可以清楚的看到,谁先执行完了谁先返回,创建的10个协程都是根据sleep的时间升序排序
当然了,swoole里面也有channel,swoole的channel用于协程直接的通讯,支持多生产者协程和多消费者协程。扩展底层自动实现了协程的切换和调度,也能够满足我们的需求,当进行多任务的时候,我们经常会需要对一些任务进行调度编排,这个时候channel就必不可少了,以下是使用channel的一段示例代码:
swoole中的channel
use Swoole\Coroutine;use Swoole\Coroutine\Channel;use function Swoole\Coroutine\run;run(function(){//新建一个通道$channel = new Channel(1);Coroutine::create(function () use ($channel) {for($i = 0; $i < 10; $i++) {Coroutine::sleep(1.0);//往通道里面塞数据$channel->push(['rand' => rand(1000, 9999), 'index' => $i]);echo "{$i}\n";}});Coroutine::create(function () use ($channel) {//等待通道出数据while(1) {$data = $channel->pop(2.0);if ($data) {var_dump($data);} else {assert($channel->errCode === SWOOLE_CHANNEL_TIMEOUT);break;}}});});
得到的部分结果如下图所示:
这段代码new了一个长度channel对象,Coroutine::create与go关键字一样,都是创建一个协程,并且在协程里循环往channel里Push数据,另一个协程负责接受数据并打印。
以上两段代码简单的演示了如何使用swoole来实现协程,是不是很简单呢?那么,用原生的PHP能否实现协程呢?
2、原生PHP通过yield关键字实现伪协程的交替多任务协作
关于生成器yield的概念与使用,可以参考一下同事们往期写的wiki,讲述的非常清楚,传送门:
【2021Q3原创-沈武斌】PHP生成器yield使用
【2021Q3原创-黎志豪】PHP yield 协程实战—"多线程"任务调度器(一)
【2021Q4原创-黎志豪】PHP yield 协程实战—"多线程"任务调度器(二)
那么为什么叫伪协程呢?
首先来看一个经典的例子:
<?phpfunction xrange($start, $end, $step = 1) {for ($i = $start; $i <= $end; $i += $step) {yield $i;}}foreach (xrange(1, 1000000) as $num) {echo $num, "\n";}
上面这个xrange()函数提供了和PHP的内建函数range()一样的功能,但是不同的是range()函数返回的是一个包含值从1到100万0的数组。而xrange()函数返回的是依次输出这些值的一个迭代器, 而不会真正以数组形式返回
这种方法的优点是显而易见的.它可以让你在处理大数据集合的时候不用一次性的加载到内存中,甚至你可以处理无限大的数据流。要从生成器认识协程, 理解它内部是如何工作是非常重要的:生成器是一种可中断的函数,在它里面的yield构成了中断点,
看回上面的例子,调用xrange(1,1000000)的时候,xrange()函数里代码其实并没有真正地运行,它只是返回了一个迭代器:
<?php$range = xrange(1, 1000000);var_dump($range); // object(Generator)#1var_dump($range instanceof Iterator); // bool(true)?>
协程的支持是在迭代生成器的基础上,,增加了可以回送数据给生成器的功能(调用者发送数据给被调用的生成器函数). 这就把生成器到调用者的单向通信转变为两者之间的双向通信。传递数据的功能是通过迭代器的send()方法实现的。下面的logger()协程是实现通信的例子:
<?phpfunction logger($fileName) {$fileHandle = fopen($fileName, 'a');while (true) {fwrite($fileHandle, yield . "\n");}}$logger = logger(__DIR__ . '/log');$logger->send('go');$logger->send('php')
这儿yield没有作为一个语句来使用,而是用作一个表达式, 即它能被演化成一个值,这个值就是调用者传递给send()方法的值。在这个例子里,yield表达式将首先被"go"替代写入Log, 然后被"php"替代写入Log。
多任务协作这个术语中的“协作”很好的说明了如何进行这种切换的:它要求当前正在运行的任务自动把控制传回给调度器,这样就可以运行其他任务了。 这与“抢占式”多任务相反, 抢占多任务是这样的:调度器可以中断运行了一段时间的任务, 不管它任务是否
愿意,协作多任务在Windows的早期版本(windows95)和Mac OS中有使用,不过它们后来都切换到使用抢先多任务了。原因也很简单:如果你依靠程序自动交出控制的话,,那么一些恶意的程序将很容易占用整个CPU,而其他任务则只能在后台看着。
yield指令提供了任务中断自身的一种方法, 然后把控制交回给任务调度器。因此协程可以运行多个其他任务,更进一步来说, yield还可以用来在任务和调度器之间进行通信。现在我们用yield来实现多任务协作,包装一个任务类和调度类,任务类负责管理任务相关的,例如执行任务,生成任务ID,调度类负责调度任务:
<?php//任务类class Task {protected $taskId;protected $coroutine;protected $sendValue = null;protected $beforeFirstYield = true;public function __construct($taskId, Generator $coroutine) {$this->taskId = $taskId;$this->coroutine = $coroutine;}public function getTaskId() {return $this->taskId;}public function setSendValue($sendValue) {$this->sendValue = $sendValue;}public function run() {//确定第一个yield的值能被正确返回if ($this->beforeFirstYield) {$this->beforeFirstYield = false;return $this->coroutine->current();} else {$retval = $this->coroutine->send($this->sendValue);$this->sendValue = null;return $retval;}}public function isFinished() {return !$this->coroutine->valid();}}class Scheduler {protected $maxTaskId = 0;protected $taskMap = []; // taskId => taskprotected $taskQueue;public function __construct() {$this->taskQueue = new SplQueue();}public function newTask(Generator $coroutine) {$tid = ++$this->maxTaskId;$task = new Task($tid, $coroutine);$this->taskMap[$tid] = $task;$this->schedule($task);return $tid;}public function schedule(Task $task) {$this->taskQueue->enqueue($task);}public function run() {while (!$this->taskQueue->isEmpty()) {$task = $this->taskQueue->dequeue();$task->run();if ($task->isFinished()) {unset($this->taskMap[$task->getTaskId()]);} else {$this->schedule($task);}}}}//任务示例function task1() {for ($i = 1; $i <= 10; ++$i) {echo "我是任务1,迭代i的值: $i".PHP_EOL;yield;}}function task2() {for ($i = 1; $i <= 5; ++$i) {echo "我是任务2,迭代i的值: $i".PHP_EOL;yield;}}//调用示例$scheduler = new Scheduler;$scheduler->newTask(task1());$scheduler->newTask(task2());$scheduler->run();
newTask()方法创建一个新任务,然后把这个任务放入任务map数组里, 接着它通过把任务放入任务队列里来实现对任务的调度,接着run()方法扫描任务队列,运行任务。如果一个任务结束了,那么它将从队列里删除, 否则它将在队列的末尾再次被调度。末尾的代码简单示例了如何调度任务。
两个任务都仅仅回显一条信息,然后使用yield把控制回传给调度器.输出结果如下:
看到结果也如我们想要的,两个任务之间是交替运行的,而在第二个任务结束后,只有第一个任务在运行。现在知道为什么说是”伪协程“了吧,实际上,使用yield来实现协程本质上还是一种同步阻塞IO模型,对于异步投递任务并没有起到什么实质性作用。那么原生PHP就没办法实现用户态非阻塞IO的协程了吗?答案是否定的,篇幅问题具体的实现可以参考国外大佬的一篇文章,写的非常详细,但是阅读需要一定时间,因为是英文的,goole的翻译略微有些尴尬,传送门:
【Cooperative multitasking using coroutines (in PHP!) 】https://www.npopov.com/2012/12/22/Cooperative-multitasking-using-coroutines-in-PHP.html
五、结语
可能有的同学会问,为什么要用PHP来实现多任务协作呢,直接用Golang不香吗。当然用Go来实现多任务协作会比原生PHP来实现要简单的多,因为Golang本身就支持我们去做这些操作。但是,业务项目需求是运行在天翼上的,而Go版的天翼(Go托管平台)还没有正式上线,所以,在自己本身感兴趣的情况下去研究这个实现的方式并且记录下来。果然,PHP还是世界上最好的语言!