您尚未登录,请登录后浏览更多内容! 登录 | 立即注册

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15618|回复: 0
打印 上一主题 下一主题

[html5] 用于实例化Worker后执行监听

[复制链接]
跳转到指定楼层
楼主
发表于 2018-12-17 21:22:08 | 只看该作者 回帖奖励 |倒序浏览 |阅读模式
  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;, M; [7 e) i% f' a9 s# x: m* [/ N
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    ) x- [1 j( ]3 z/ S5 Y& J

  3. 7 |( w& n- R' O
  4. $worker = new Worker();% b% {' A0 u. l
  5. // 4个进程! F& l" D5 P( V" X1 U  Y8 O" A
  6. $worker->count = 4;
    1 d$ I! I1 K4 T5 a$ V" F
  7. // 每个进程启动后在当前进程新增一个Worker监听# c3 S% ]1 y; @- E
  8. $worker->onWorkerStart = function($worker)
      C' A/ J' `& S+ F
  9. {
      w1 _* p; p" a: U
  10.     /**
    , l/ t: c3 r6 U0 C
  11.      * 4个进程启动的时候都创建2016端口的Worker
    4 r6 R+ @1 ^& c9 }
  12.      * 当执行到worker->listen()时会报Address already in use错误
    : C4 [8 V  ]- `" V% Q
  13.      * 如果worker->count=1则不会报错9 t* c# Q: t4 K
  14.      */
      h8 f' e. Q7 ^
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');
    6 g8 q2 h) I6 s: {8 M) {8 |: N
  16.     $inner_worker->onMessage = 'on_message';; N4 J' J$ E) c4 t1 Y
  17.     // 执行监听。这里会报Address already in use错误; B/ f( h* R8 Q+ K# i
  18.     $inner_worker->listen();
    & p6 d' m6 g! q- n+ P, f( @
  19. };3 ]2 c" ~4 k- l, O5 b% d- S
  20. 0 n' e; A2 F$ _5 g0 j
  21. $worker->onMessage = 'on_message';
    % ^% L( ]$ \4 n/ F6 {: |  ?4 @
  22. ( b+ {: K( J5 U5 K: m9 x7 ^# d( }
  23. function on_message($connection, $data)
    ' L! s4 \2 ~7 L, A
  24. {% B, {2 M, p# X2 P1 K! V
  25.     $connection->send("hello\n");
    # o3 G0 E/ P+ ~* X" m
  26. }
    9 f, d2 p8 t. A5 q/ ~

  27. , i) Q: n; M* I
  28. // 运行worker/ ^. o. ]8 @. V
  29. Worker::runAll();1 ~2 M" h6 \# b$ @9 k% H  b4 n
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:1 p+ ?. a  c# d/ N' K
  31. 1 f, s# B0 p' y( L
  32. use Workerman\Worker;
    " K. H5 J2 H5 C8 d, k; G' H5 V9 u% t
  33. require_once './Workerman/Autoloader.php';
    1 D* S$ J- q8 b# w0 O$ s# o0 s

  34. " K, ~  x: M9 y+ p$ H) g( ]
  35. $worker = new Worker('text://0.0.0.0:2015');: `) L& \* _! Z' J  v* P. [5 g6 Z
  36. // 4个进程; c% @' c7 v' i  d4 _( [4 l
  37. $worker->count = 4;+ P& H8 \/ l6 x! L
  38. // 每个进程启动后在当前进程新增一个Worker监听" K  O! a5 O7 {
  39. $worker->onWorkerStart = function($worker)
    ! _# i7 y: w: X$ S$ c; [; b: Z
  40. {0 v: }5 @: C2 o3 h/ Q4 Z7 f
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');  m2 G# U$ {4 V) v5 b
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)5 I, r+ u& D+ f8 O9 }6 F$ i
  43.     $inner_worker->reusePort = true;$ \5 J+ b6 `% R8 C6 W( y% v
  44.     $inner_worker->onMessage = 'on_message';) }/ K: d: ~- P
  45.     // 执行监听。正常监听不会报错
    ; v& h0 Y, N# o9 ~( V8 `
  46.     $inner_worker->listen();
    , q. r, `8 F4 u5 B/ u0 V
  47. };
    3 T) M' I0 m, |; |, n5 A, Z7 H+ ^

  48. 4 F/ g6 h6 O4 ]
  49. $worker->onMessage = 'on_message';
    1 e/ [& z9 u( L
  50. 8 E- Y# d( e1 ?8 i
  51. function on_message($connection, $data)
      ^7 W% T/ g' l
  52. {
    ' k; g. f2 }# U' q
  53.     $connection->send("hello\n");
    / p) Y7 m$ }4 M. W3 g
  54. }
    & e# Z9 F6 U) t7 r* j% g* j5 K
  55. ! b9 p8 H. E; B  b% ?
  56. // 运行worker  S+ v7 Q) J/ ^* O$ _
  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
    2 @. Y& V  ^, C" x& E6 |0 L
  2. use Workerman\Worker;
    - o# N0 J6 M! n0 X5 n
  3. require_once './Workerman/Autoloader.php';# v& S' S# s6 E: h& ]1 r& @0 c
  4. // 初始化一个worker容器,监听1234端口+ h+ o% r2 B6 O9 S( h
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    2 g3 U3 P6 J: ~

  6. : W; H9 B  f. c( \  e
  7. /*0 D5 H  `5 H' v! B+ T
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    - D% U, n7 W8 E# T2 N8 \( D4 `
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true), l, m0 t; a- C7 R0 @, j/ q
  10. */
      s3 d3 R. y# y' P# N" J
  11. $worker->count = 1;: N- ?- `+ E. Y
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口. x% u- u3 n, u$ K
  13. $worker->onWorkerStart = function($worker)! I  E9 B3 k4 \. g7 y$ U8 M4 C
  14. {
    # Z8 g2 `! G2 n1 s
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符7 ~( Z) N" T6 b+ _
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');  M  T8 s2 u; s0 m# R
  17.     $inner_text_worker->onMessage = function($connection, $buffer)) O( ^( ^; M+ R6 ^
  18.     {6 ^5 F' e, o$ N: ~/ X& x+ D2 u& t
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据
    % X& @5 V6 U2 W! y
  20.         $data = json_decode($buffer, true);) H$ D- M' R6 E7 O$ E* X- ?: J6 B
  21.         $uid = $data['uid'];
    2 E6 R4 M9 W5 L' c6 O+ k1 s" g
  22.         // 通过workerman,向uid的页面推送数据- v* S3 o( c1 A$ x
  23.         $ret = sendMessageByUid($uid, $buffer);
    - }$ R, S9 b& n. F5 ~; L2 z- @
  24.         // 返回推送结果
    5 U, Z6 x4 i4 l# D: N
  25.         $connection->send($ret ? 'ok' : 'fail');! F  O! N/ {5 p" W; {$ z* N" q
  26.     };$ Y7 ~5 Y% p6 d% j
  27.     // ## 执行监听 ##
    ' {. ^& {' q! D  n
  28.     $inner_text_worker->listen();/ {4 O/ I9 z5 D- [; l/ j
  29. };; U/ ~! S$ w8 a2 b: z
  30. // 新增加一个属性,用来保存uid到connection的映射
    7 d" x, A& V" [" |& A. L* M) g
  31. $worker->uidConnections = array();# v* }4 q3 L- v9 S- j
  32. // 当有客户端发来消息时执行的回调函数0 t  \7 ~* r# }5 K
  33. $worker->onMessage = function($connection, $data)
    4 h: J& G. l# C4 W0 [3 i2 t' D' Z
  34. {+ V3 l1 m6 H1 Q$ ^- l( r0 k8 U6 I3 n
  35.     global $worker;
    ( O- I8 V1 Z; L; K
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    ( l+ [1 W' K* Y/ f# A, E
  37.     if(!isset($connection->uid))2 V- t2 h" O; t- U' t1 t/ c, _+ N) \
  38.     {9 L/ f# v* k6 `- U& l0 `% \+ L
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)2 s. s6 f. s, ?
  40.        $connection->uid = $data;
    $ Y* q# c/ q& f! I- }7 S2 Q
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,; T2 r% O. J3 ]1 V1 w6 j
  42.         * 实现针对特定uid推送数据( v( M' M$ ]% ?- V$ z6 N
  43.         */; _5 N! P* O/ G: R+ `8 |$ @
  44.        $worker->uidConnections[$connection->uid] = $connection;
    2 S) u: ?& f8 s. c, S5 J% S* V
  45.        return;
    - r; x) F( B( {% }6 f' W" q
  46.     }! D' F& k  E. ^8 j4 T+ c3 j( ?: V
  47. };' U6 I4 L( K+ I7 X- b1 ?$ {( ~
  48. # `& p" O6 n; X- U" S& Q  j$ D
  49. // 当有客户端连接断开时; C, C# U# d" _9 S
  50. $worker->onClose = function($connection): G2 c% q! ]# E$ H/ j8 u
  51. {
    " L, O) x$ T. N6 L5 n
  52.     global $worker;
      V* T: k, B, y0 @3 Q' Q
  53.     if(isset($connection->uid))7 k' f, Y0 M( q4 Q! s" X
  54.     {. N9 y, L6 c/ m2 H: _/ T
  55.         // 连接断开时删除映射/ [, z& F5 E2 u" {2 Q4 t& k2 ]
  56.         unset($worker->uidConnections[$connection->uid]);
    + x6 n3 c+ f! D& r
  57.     }
    # n1 u$ \1 U2 L9 C( B9 v9 S
  58. };0 }1 b# P4 A% `, j( m5 ?  ^+ M

  59. 9 A/ h7 X7 r& @/ ^
  60. // 向所有验证的用户推送数据1 r' O$ U) ?  B$ j+ @7 f
  61. function broadcast($message)
    ( T5 U4 x- j8 i! i8 O
  62. {
    5 z# Y2 }5 H4 _- T& K2 u
  63.    global $worker;* A7 T. ?& D2 m5 x' n9 l: ~% U
  64.    foreach($worker->uidConnections as $connection)+ v: {! D7 _+ M
  65.    {1 D" F# f/ J0 ]  j1 }  u1 K
  66.         $connection->send($message);$ p4 p( W: W6 h0 Z, s
  67.    }/ {3 [' ~& ^/ G- I  j$ D8 q- L
  68. }
    $ H6 {2 ?9 R3 @* O5 |
  69. 4 E8 ~1 S- K' E2 O* K) v( `% V; z
  70. // 针对uid推送数据! l) H7 ]4 J, b/ j8 [/ i
  71. function sendMessageByUid($uid, $message)
    # m& |. Y+ I' G5 n% v! e3 I# K5 Z
  72. {
    " o7 |# j$ h0 ~7 ?' K
  73.     global $worker;
    6 h* a4 d( j) [: M
  74.     if(isset($worker->uidConnections[$uid]))
    2 Q( D- y5 V$ S  U9 ?
  75.     {
    ) Z! V# H3 c, w
  76.         $connection = $worker->uidConnections[$uid];
    . i3 \2 _+ J) y. _6 b
  77.         $connection->send($message);) w5 X, H+ C- E2 a; I
  78.         return true;# n, h  Z* c  @' k; M0 ~7 o: d! {. n
  79.     }
      T( {: j! A; `
  80.     return false;3 C8 Y% {! j8 V2 u
  81. }
      v' x' G' d  W: t; ]# }

  82. ) }) H) E% K# u7 V' k' |
  83. // 运行所有的worker
    , A2 e2 V' s: g) ^1 z' s
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');" n6 K7 i5 [! d6 c0 w- w( L% G5 t
  2. ws.onopen = function(){
    6 T/ Y5 v+ G5 l+ ~& q2 G, ?
  3.     var uid = 'uid1';
    1 T$ I- j( Y1 `$ k2 |3 F1 {
  4.     ws.send(uid);, d: \0 w  o" T: |& T8 P3 M
  5. };
    : I) u; S5 |2 U6 D3 t
  6. ws.onmessage = function(e){
    . O, a4 v" P. o) e4 k) Y9 c% z
  7.     alert(e.data);# ~5 W+ V7 B; ^9 m3 q, f: ?( F& p
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    - S8 U' |1 s' S% _; ~7 l% w, O
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);$ t. e* ]4 E. x# n% R& F
  3. // 推送的数据,包含uid字段,表示是给这个uid推送
    $ h# v' k& E2 ~( N+ f: z
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');
    # q- ]% _$ W2 \/ d/ \% ~
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符# V! R* U* ^& V4 b, G7 S
  6. fwrite($client, json_encode($data)."\n");
    - Q4 n4 O8 S8 a/ o$ h
  7. // 读取推送结果
    + s: |+ x- v/ D
  8. echo fread($client, 8192);
复制代码

4 Q9 Q- q& X- [1 m% y
. d- g  z2 c) c3 M' m+ S- I
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 13:54 , Processed in 0.057598 second(s), 21 queries .

Copyright © 2001-2026 Powered by cncml! X3.2. Theme By cncml!