<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->serv = new swoole_server("0.0.0.0", 9501);
6 $this->serv->set(array(
7 'worker_num' => 4,
8 'daemonize' => false,
9 'task_worker_num' => 8
10 ));
11 $this->serv->on('Start', array($this, 'onStart'));
12 $this->serv->on('Connect', array($this, 'onConnect'));
13 $this->serv->on('Receive', array($this, 'onReceive'));
14 $this->serv->on('Close', array($this, 'onClose'));
15 $this->serv->on('Task', array($this, 'onTask'));
16 $this->serv->on('Finish', array($this, 'onFinish'));
17 $this->serv->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' => $fd
25 );
26 $serv->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 < 5 ; $i ++ ) {
34 sleep(1);
35 echo "Task {$task_id} Handle {$i} times...\n";
36 }
37 $fd = json_decode( $data , true )['fd'];
38 $serv->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->serv = new swoole_server("0.0.0.0", 9501);
6 $this->serv->set(array(
7 'worker_num' => 4,
8 'daemonize' => false,
9 'task_worker_num' => 8 // task进程数量 即为维持的MySQL连接的数量
10 ));
11 $this->serv->on('Start', array($this, 'onStart'));
12 $this->serv->on('Connect', array($this, 'onConnect'));
13 $this->serv->on('Receive', array($this, 'onReceive'));
14 $this->serv->on('Close', array($this, 'onClose'));
15 $this->serv->on('Task', array($this, 'onTask'));
16 $this->serv->on('Finish', array($this, 'onFinish'));
17 $this->serv->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' => $data, // 接收客户端发送的 sql
25 'fd' => $fd
26 );
27 $serv->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->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->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->serv->send($result['fd'], json_encode($result['data']) . "\n");
66 } else {
67 $this->serv->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