37DATA

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;    //主进程的pid    protected $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();

执行的结果如下:

Image

根据打印的顺序,我们可以知道,最开始执行完成的是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; }); }});

运行两次结果如下:

ImageImage

我们可以清楚的看到,谁先执行完了谁先返回,创建的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; } } });});

得到的部分结果如下图所示:

ImageImageImage

这段代码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 => task protected $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把控制回传给调度器.输出结果如下:

Image

看到结果也如我们想要的,两个任务之间是交替运行的,而在第二个任务结束后,只有第一个任务在运行。现在知道为什么说是”伪协程“了吧,实际上,使用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还是世界上最好的语言!

Image