PHP 使用 Swoole

<p>在一般的 Server 程序中都会有一些耗时的任务,比如:发送邮件、聊天服务器发送广播等。如果我们采用同步阻塞的防水去执行这些任务,那么这肯定会非常的慢。</p> <p>Swoole 的 TaskWorker 进程池可以用来执行一些异步的任务,而且不会影响接下来的任务,很适合处理以上场景。</p> <p>那么什么是异步任务呢?</p> <p>可以从下面的图示中来简单了解一下。(来源于网络,侵删)<br></p>

<p>我们上一个 Swoole 的文章介绍了如何创建一个简单的服务器,并且知道了几个核心的回调函数的使用方法。</p> <p>要实现上述的异步处理,只需要增加两个事件回调即可:onTask 和 onFinish, 这两个回调函数分别用于执行 Task 任务和处理 Task 任务的返回结果。另外还需要在 set 方法中设置 task 进程数量。 </p> <p>使用示例:</p>
1class Server 2{ 3 private $serv; 4 public function __construct() { 5 $this-&gt;serv = new swoole_server("0.0.0.0", 9501); 6 $this-&gt;serv-&gt;set(array( 7 'worker_num' =&gt; 4, 8 'daemonize' =&gt; false, 9 'task_worker_num' =&gt; 8 10 )); 11 $this-&gt;serv-&gt;on('Start', array($this, 'onStart')); 12 $this-&gt;serv-&gt;on('Connect', array($this, 'onConnect')); 13 $this-&gt;serv-&gt;on('Receive', array($this, 'onReceive')); 14 $this-&gt;serv-&gt;on('Close', array($this, 'onClose')); 15 $this-&gt;serv-&gt;on('Task', array($this, 'onTask')); 16 $this-&gt;serv-&gt;on('Finish', array($this, 'onFinish')); 17 $this-&gt;serv-&gt;start(); 18 } 19 20 public function onReceive( swoole_server $serv, $fd, $from_id, $data ) { 21 echo "Get Message From Client {$fd}:{$data}\n"; 22 // 发送任务到Task进程 23 $param = array( 24 'fd' =&gt; $fd 25 ); 26 $serv-&gt;task( json_encode( $param ) ); 27 echo "继续处理之后的逻辑\n"; 28 } 29 30 public function onTask($serv, $task_id, $from_id, $data) { 31 echo "This Task {$task_id} from Worker {$from_id}\n"; 32 echo "Data: {$data}\n"; 33 for($i = 0 ; $i &lt; 5 ; $i ++ ) { 34 sleep(1); 35 echo "Task {$task_id} Handle {$i} times...\n"; 36 } 37 $fd = json_decode( $data , true )['fd']; 38 $serv-&gt;send( $fd , "Data in Task {$task_id}"); 39 return "Task {$task_id}'s result"; 40 } 41 public function onFinish($serv,$task_id, $data) { 42 echo "Task {$task_id} finish\n"; 43 echo "Result: {$data}\n"; 44 } 45 public function onStart( $serv ) { 46 echo "Server Start\n"; 47 } 48 public function onConnect( $serv, $fd, $from_id ) { 49 echo "Client {$fd} connect\n"; 50 } 51 public function onClose( $serv, $fd, $from_id ) { 52 echo "Client {$fd} close connection\n"; 53 } 54} 55$server = new Server(); 56 57
<p>通过上述示例可以看到,发起一个异步任务只需要调用 swoole\_server 的 task 方法就可以。发送之后会触发 onTask 回调,可以通过 $task\_id 和 $from\_id 处理不同进程的不同任务。最后可以通过 return 一个字符串来将执行结果返回给 Worker 进程,Worker 进程通过 onFinish 回调来处理结果。</p> <p>那么基于上述代码就可以实现异步操作 mysql。异步操作 mysql 较适合以下场景:</p> <ul> <li>并发的读写操作</li> <li>没有时序上的严格关系</li> <li>不影响主线程逻辑</li> </ul> <p>好处:</p> <ul> <li>提高并发</li> <li>降低 IO 消耗</li> </ul> <p>数据库的压力主要在于 mysql 维持的连接数,如果存在 1000 个并发,那么 mysql 就需要建立对应数量的连接。而采用长连接的方式,mysql 的连接一直维持在进程中,减少了创建连接的损耗。可以通过 swoole 开启多个 task 进程,每一个进程内维持一个mysql 长连接,那么这样子也可以引申出来 mysql 连接池技术。还需要注意的是,mysql 服务器如果检测到长时间没有没有查询,则会断开连接回收资源,所以要有断线重连的机制。</p> <p>以下是一个简单的异步操作 mysql 的示例:</p> <p>还是以上的代码,我们只需要修改 onReceive、onTask、onFinish 三个函数。</p>
1class Server 2{ 3 private $serv; 4 public function __construct() { 5 $this-&gt;serv = new swoole_server("0.0.0.0", 9501); 6 $this-&gt;serv-&gt;set(array( 7 'worker_num' =&gt; 4, 8 'daemonize' =&gt; false, 9 'task_worker_num' =&gt; 8 // task进程数量 即为维持的MySQL连接的数量 10 )); 11 $this-&gt;serv-&gt;on('Start', array($this, 'onStart')); 12 $this-&gt;serv-&gt;on('Connect', array($this, 'onConnect')); 13 $this-&gt;serv-&gt;on('Receive', array($this, 'onReceive')); 14 $this-&gt;serv-&gt;on('Close', array($this, 'onClose')); 15 $this-&gt;serv-&gt;on('Task', array($this, 'onTask')); 16 $this-&gt;serv-&gt;on('Finish', array($this, 'onFinish')); 17 $this-&gt;serv-&gt;start(); 18 } 19 20 public function onReceive( swoole_server $serv, $fd, $from_id, $data ) { 21 echo "收到数据". $data . PHP_EOL; 22 // 发送任务到Task进程 23 $param = array( 24 'sql' =&gt; $data, // 接收客户端发送的 sql 25 'fd' =&gt; $fd 26 ); 27 $serv-&gt;task( json_encode( $param ) ); // 向 task 投递任务 28 echo "继续处理之后的逻辑\n"; 29 } 30 31 public function onTask($serv, $task_id, $from_id, $data) { 32 echo "This Task {$task_id} from Worker {$from_id}\n"; 33 echo "recv SQL: {$data['sql']}\n"; 34 static $link = null; 35 $sql = $data['sql']; 36 $fd = $data['fd']; 37 HELL: 38 if ($link == null) { 39 $link = @mysqli_connect("127.0.0.1", "root", "root", "test"); 40 } 41 $result = $link-&gt;query($sql); 42 if (!$result) { //如果查询失败 43 if(in_array(mysqli_errno($link), [2013, 2006])){ 44 //错误码为2013,或者2006,则重连数据库,重新执行sql 45 $link = null; 46 goto HELL; 47 } 48 } 49 if(preg_match("/^select/i", $sql)){//如果是select操作,就返回关联数组 50 $data = array(); 51 while ($fetchResult = mysqli_fetch_assoc($result) ){ 52 $data['data'][] = $fetchResult; 53 } 54 }else{//否则直接返回结果 55 $data['data'] = $result; 56 } 57 $data['status'] = "OK"; 58 $data['fd'] = $fd; 59 $serv-&gt;finish(json_encode($data)); 60 } 61 public function onFinish($serv, $task_id, $data) { 62 echo "Task {$task_id} finish\n"; 63 $result = json_decode($result, true); 64 if ($result['status'] == 'OK') { 65 $this-&gt;serv-&gt;send($result['fd'], json_encode($result['data']) . "\n"); 66 } else { 67 $this-&gt;serv-&gt;send($result['fd'], $result); 68 } 69 } 70 public function onStart( $serv ) { 71 echo "Server Start\n"; 72 } 73 public function onConnect( $serv, $fd, $from_id ) { 74 echo "Client {$fd} connect\n"; 75 } 76 public function onClose( $serv, $fd, $from_id ) { 77 echo "Client {$fd} close connection\n"; 78 } 79} 80$server = new Server(); 81
<p>以上代码在 onReceive 时直接接收一条 sql,之后直接发送到 Task 任务中。这个时候下一步的流程紧接着输出,这里也就体现出了异步。然后 onTask 和 onFinish 分别用来向数据库发送 sql,处理 task 执行结果。</p> <p>参考链接:</p> <p><a href="https://wiki.swoole.com" rel="nofollow noreferrer">https://wiki.swoole.com</a><br><a href="http://rango.swoole.com/archives/265" rel="nofollow noreferrer">http://rango.swoole.com/archi...</a></p>

原文地址:https://segmentfault.com/a/1190000016706048

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )