cncml手绘网

标题: 用于实例化Worker后执行监听 [打印本页]

作者: admin    时间: 2018-12-17 21:22
标题: 用于实例化Worker后执行监听
  1. void Worker::listen(void)
复制代码
用于实例化Worker后执行监听。
此方法主要用于在Worker进程启动后动态创建新的Worker实例,能够实现同一个进程监听多个端口,支持多种协议。需要注意的是用这种方法只是在当前进程增加监听,并不会动态创建新的进程,也不会触发onWorkerStart方法。
例如一个http Worker启动后实例化一个websocket Worker,那么这个进程即能通过http协议访问,又能通过websocket协议访问。由于websocket Worker和http Worker在同一个进程中,所以它们可以访问共同的内存变量,共享所有socket连接。可以做到接收http请求,然后操作websocket客户端完成向客户端推送数据类似的效果。
注意:
如果PHP版本<=7.0,则不支持在多个子进程中实例化相同端口的Worker。例如A进程创建了监听2016端口的Worker,那么B进程就不能再创建监听2016端口的Worker,否则会报Address already in use错误。例如下面的代码是无法运行的。
  1. use Workerman\Worker;  f+ r6 P9 e3 t/ O; N
  2. require_once __DIR__ . '/Workerman/Autoloader.php';, X- e# U, {  Z: h1 s

  3. % F! F& h  M/ i: p7 f; Z" k
  4. $worker = new Worker();
    * X. t# P" f8 z4 F0 F9 t0 l
  5. // 4个进程: p" X* I! C& d! Y% p: y! G
  6. $worker->count = 4;3 ?6 T: D' E4 `7 c/ h7 `+ h
  7. // 每个进程启动后在当前进程新增一个Worker监听
    : K' V- e( P7 L4 `4 t4 P  i0 ]
  8. $worker->onWorkerStart = function($worker)# j9 o; h+ [* t1 ~0 P. M0 \# f
  9. {2 I1 ^$ `8 [! J+ J
  10.     /**4 E) W2 e/ ^1 ]; k/ Z: H; l
  11.      * 4个进程启动的时候都创建2016端口的Worker. V2 c# A" u% c
  12.      * 当执行到worker->listen()时会报Address already in use错误
    & E5 v9 G* N- p; D: j3 l- j+ }
  13.      * 如果worker->count=1则不会报错" l; l) X4 I  {1 _' z
  14.      */. @4 S/ f* [# {0 {# i1 k
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');( ]9 m. i9 ?" ?3 i0 W  y
  16.     $inner_worker->onMessage = 'on_message';
    ' [$ i; H2 G* M/ Z# i
  17.     // 执行监听。这里会报Address already in use错误7 T9 H0 X0 t$ z- j% u# c
  18.     $inner_worker->listen();: @8 _; M% Z9 ?5 u& a7 w& e# {
  19. };/ ^3 U* P' P, \+ t- C* l
  20. " M; ?  L- y  k
  21. $worker->onMessage = 'on_message';6 {2 `# k4 F. M

  22. ( A' N: d! ~( D& r4 c0 f4 {3 R
  23. function on_message($connection, $data)
    % f+ m5 A7 t* h
  24. {' A, E% ?! U5 ]0 F: |9 g
  25.     $connection->send("hello\n");
    8 r! o4 B9 ~, }( D" S. z# U! u% D
  26. }5 T7 n, q! C3 _4 n
  27. 7 P/ F" B$ Q/ _2 y! \( _: @, O
  28. // 运行worker4 l+ I  @( ]$ x" s( p" n$ C
  29. Worker::runAll();4 G% h# Z. d6 G& }) j  a& H8 f
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:  A" A( W3 C8 X1 ]' ~" B4 Z7 K
  31. 1 u& T. ]" l2 o0 A4 K$ `. J; a
  32. use Workerman\Worker;, N! U4 D1 t" N7 C# s
  33. require_once './Workerman/Autoloader.php';
    3 a8 D/ D" J. n: w

  34. 7 v! l( T2 }: w/ a: E
  35. $worker = new Worker('text://0.0.0.0:2015');
    1 c& Z9 F, L- D9 D5 R
  36. // 4个进程
    8 W4 |0 ~2 z! p6 H8 @' D5 w$ T
  37. $worker->count = 4;
      [# W( M  K* v, }2 e6 y
  38. // 每个进程启动后在当前进程新增一个Worker监听/ t4 L7 k) ^0 J9 g& G& b3 C
  39. $worker->onWorkerStart = function($worker)# t5 ~4 E1 p. F4 g+ W0 c, G
  40. {; [9 D" p, u* Q" U$ x9 s
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');3 z8 P' {. ]" \5 A8 t: ]
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    ! k1 D# Q" x3 e) _9 U/ _8 b
  43.     $inner_worker->reusePort = true;2 R/ X5 s$ [3 }1 B5 c2 V# H" u: y' ?, f
  44.     $inner_worker->onMessage = 'on_message';
    ' ], y/ ^! F' L$ `4 I0 V$ A
  45.     // 执行监听。正常监听不会报错  y1 b: a+ R) K
  46.     $inner_worker->listen();0 Y+ b/ B! X/ p
  47. };
    / b, B. @/ G! v
  48. + p* K  [' J( q* D$ i4 k
  49. $worker->onMessage = 'on_message';% z* \2 f: G1 l8 {9 q  ]8 z$ V

  50. . X( {. H, N5 T" }& f
  51. function on_message($connection, $data)0 _0 ^4 ^0 h4 v7 ^
  52. {
    9 R( q% c  c# M' {3 z' M
  53.     $connection->send("hello\n");
    - Q4 u0 D0 V! T1 n9 {+ w/ J! i
  54. }0 z2 ]$ H4 Z6 F9 {* S( `

  55. # H& g  C8 Z. J6 B, z
  56. // 运行worker
    . p' @; u  q7 ?
  57. Worker::runAll();
复制代码
示例 php后端及时推送消息给客户端
原理:
1、建立一个websocket Worker,用来维持客户端长连接
2、websocket Worker内部建立一个text Worker
3、websocket Worker 与 text Worker是同一个进程,可以方便的共享客户端连接
4、某个独立的php后台系统通过text协议与text Worker通讯
5、text Worker操作websocket连接完成数据推送
代码及步骤
push.php
  1. <?php
    ' p- r# O4 _# K
  2. use Workerman\Worker;
    + n4 i! W" _6 Z* C' f
  3. require_once './Workerman/Autoloader.php';
    6 t1 Q, h+ |0 K( \( G* \+ y3 g
  4. // 初始化一个worker容器,监听1234端口) _; m4 J8 w7 y
  5. $worker = new Worker('websocket://0.0.0.0:1234');- S+ S9 z7 j$ }5 \4 f$ B9 [
  6. 5 {6 z* D7 i. J# t& ~- R; N
  7. /*
    4 F1 ^/ Q- O* x5 H4 q) M
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误& F9 w, I7 Q; d* u9 Z$ l' v4 W
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)5 u" Q! B9 ?2 e& q% G# Q, x
  10. */
    5 k6 x# c# I  s% [8 \: y# h4 l3 s
  11. $worker->count = 1;$ ?/ ^" W$ E% Q. }! j# s" f* K" _5 v
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口6 P& S; p4 u3 d7 m$ N
  13. $worker->onWorkerStart = function($worker)
    : k! i/ ]- H! N9 C0 W1 `3 G, N2 v
  14. {* A1 S/ C/ w. ^2 }0 y. @% Y
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    6 f% A; b% I0 M+ m
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');$ t# o* y* ~3 ]7 ]) I
  17.     $inner_text_worker->onMessage = function($connection, $buffer)
    * X/ [# o0 O# s# R3 F( `
  18.     {# E6 c4 Z8 N2 G" l4 b# q
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据
    7 X# q( k/ u3 h# C
  20.         $data = json_decode($buffer, true);- N$ J- c+ L* ^! X- e
  21.         $uid = $data['uid'];7 y5 T/ Z; [) x
  22.         // 通过workerman,向uid的页面推送数据
    4 ?1 a( [& _3 z9 F5 k/ V! [
  23.         $ret = sendMessageByUid($uid, $buffer);# n6 n5 @0 \0 ~& v
  24.         // 返回推送结果
    0 e$ {% l% q! K1 H
  25.         $connection->send($ret ? 'ok' : 'fail');
    , Z' H: W. o' m( M3 g$ q; t+ p
  26.     };
    & e8 b3 |5 k& p9 Q
  27.     // ## 执行监听 ##
    9 j: H- ]+ P" y* @& b/ W5 b! |
  28.     $inner_text_worker->listen();
    + a) n: E/ ^9 X
  29. };8 N: x( z' [+ P
  30. // 新增加一个属性,用来保存uid到connection的映射
    / M) ~6 b6 _& [
  31. $worker->uidConnections = array();! w; |: P1 f/ O5 v( v
  32. // 当有客户端发来消息时执行的回调函数. J' d; `- {  X( j
  33. $worker->onMessage = function($connection, $data)2 U: B. S0 E5 s" x
  34. {: J- Y0 X7 ^* B6 T: w
  35.     global $worker;
    $ W4 t) ~8 r8 y( g& Y# |1 t
  36.     // 判断当前客户端是否已经验证,既是否设置了uid+ ~( L0 r' `2 R% x
  37.     if(!isset($connection->uid))% W! o; u# Y! E
  38.     {
    / s# U. U4 D0 w, ?. v8 K7 u2 T2 w0 X
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
    # W" p/ g5 K( x) y2 T' N: d8 p) j
  40.        $connection->uid = $data;
      ?& n. @8 n+ R
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,+ A' V2 V/ k0 S1 o
  42.         * 实现针对特定uid推送数据
    " P" R+ {5 ~! u) _; s1 I+ e$ Y
  43.         */" C% ^) s+ u$ t
  44.        $worker->uidConnections[$connection->uid] = $connection;+ l# k5 E/ z  g. B) _
  45.        return;
    9 P* e3 T- O0 N1 m
  46.     }6 H0 n6 x9 E6 `, c, J
  47. };
    % h) i- C% s5 x! O* d/ a

  48. 1 m- g' a2 ?- C
  49. // 当有客户端连接断开时- k- m/ Z) `7 ?# c, M
  50. $worker->onClose = function($connection)
    ( l  V# S; A/ q5 N9 r" W' ^
  51. {
    3 j! @2 _' N0 _8 ]  q( N
  52.     global $worker;8 E( \. v. W: e1 ^+ X+ A5 N
  53.     if(isset($connection->uid))
    2 }) v5 M+ V* c. Y6 x
  54.     {
      T) t. K! m0 s/ \4 C& Y
  55.         // 连接断开时删除映射
    $ R6 l, m! O, C2 k
  56.         unset($worker->uidConnections[$connection->uid]);+ H" W+ E0 a" f* |, u2 |2 R0 m
  57.     }
    + r8 @6 d+ N8 K) p" `& A, k3 K1 F
  58. };
    3 m& z6 v8 R% ?" @4 z. e

  59. . T! Z7 i! _& \. Q$ `
  60. // 向所有验证的用户推送数据
    ( F9 v7 [' |+ v& Z9 [# C
  61. function broadcast($message)9 i7 c/ Y, t' D. C% d) c
  62. {
    9 [/ C- i; |0 {+ \' Q9 r
  63.    global $worker;
    + g; c& J& k1 Q4 Z* Z: {- n' i& i
  64.    foreach($worker->uidConnections as $connection)/ U; L7 _/ |1 v4 S* G+ Y5 {0 c
  65.    {* u3 P+ l' f9 u3 N4 U# n
  66.         $connection->send($message);  W6 L& S2 n& ^$ j
  67.    }8 @6 N& C: I& }$ a4 e3 {
  68. }9 ^% @+ s' I3 v( a/ P9 \

  69. 7 ^  @) S; u% U9 J
  70. // 针对uid推送数据
    " P+ K/ ?! K+ i2 s0 |7 ^4 J0 O7 W
  71. function sendMessageByUid($uid, $message)+ O/ j' W3 G, A1 m
  72. {
    5 f( D. i' d3 P, h: N
  73.     global $worker;! G: \  R, h) L" L1 _0 I
  74.     if(isset($worker->uidConnections[$uid]))
    ! y) O0 p; u) X3 {
  75.     {+ r- v% a7 A: ]' u
  76.         $connection = $worker->uidConnections[$uid];2 d: f% v; Y) V: C9 |
  77.         $connection->send($message);, ]/ e/ l% G! Z) j5 n0 j
  78.         return true;% L; ^2 X( q. \- F$ @
  79.     }
      Z9 u& R6 W# Z
  80.     return false;" a2 X+ [% i+ Q3 l' J3 p
  81. }* w/ f8 B; m1 D" F  m

  82. , M  q! Q7 M5 n- O- G) P
  83. // 运行所有的worker
    5 _6 U, t8 M8 y0 h$ Q8 T
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');9 v% v$ I/ g7 l2 e0 k
  2. ws.onopen = function(){
      w6 h7 l, b$ w' x7 F9 I
  3.     var uid = 'uid1';
    # [0 {5 N" \; m8 K" k
  4.     ws.send(uid);4 |4 v7 @% V. p0 F5 t) u
  5. };
    7 Q- g( p9 k; H- s( @! ]/ Y0 b
  6. ws.onmessage = function(e){- U# p! s! Y1 C. r, s
  7.     alert(e.data);" T% S! N4 s2 X( {
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    2 K, L! `/ e# n) C; V2 Z8 h7 c, p
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);  X, a, D8 f5 j4 A! V4 L
  3. // 推送的数据,包含uid字段,表示是给这个uid推送8 q- X2 S5 d0 a9 @( q; y! ^
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');, k2 X. b$ R- Z0 z8 L
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符; Z% D: S, q  S% G& B% b
  6. fwrite($client, json_encode($data)."\n");/ b7 }- L4 k" i
  7. // 读取推送结果
    / M- X4 e2 l- V7 c* i' l
  8. echo fread($client, 8192);
复制代码
9 c* O1 t7 N* {7 U

  e4 T7 b. }0 P+ u# w$ i




欢迎光临 cncml手绘网 (http://bbs.cncml.com/) Powered by Discuz! X3.2