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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15621|回复: 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;+ g3 }( H  [( U4 S9 J, C
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    / W* C* L, P( q1 E" i

  3. $ a' k$ e, n: W+ u% [# i
  4. $worker = new Worker();, u9 C) r, C4 {% x7 g7 h1 \
  5. // 4个进程  S( R7 q& B4 R
  6. $worker->count = 4;, l% }" d$ ~: B( d0 l
  7. // 每个进程启动后在当前进程新增一个Worker监听
    , U" r5 Q2 E0 g) R: `  O
  8. $worker->onWorkerStart = function($worker)- E- f% }( h4 A( t4 i$ I
  9. {
    / E2 p1 B5 G5 Z4 ]
  10.     /**3 {% ?$ Z, i: f9 Q" A  a
  11.      * 4个进程启动的时候都创建2016端口的Worker
    0 }6 [) v* B4 b2 _/ E
  12.      * 当执行到worker->listen()时会报Address already in use错误% t6 X# [% Z) b7 Y7 Z
  13.      * 如果worker->count=1则不会报错; ?0 k5 ]' C$ f: y8 F' U8 q% ?6 m
  14.      */* \: ?& [: l: o# \% ?0 A
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');5 C! m7 v) i% ]" P7 j
  16.     $inner_worker->onMessage = 'on_message';' D/ t$ J! J( B# M2 Y
  17.     // 执行监听。这里会报Address already in use错误
    0 _  r6 c5 y% C- B. V7 ]/ z! o! z
  18.     $inner_worker->listen();$ r0 T: a& d7 O/ u) ?2 J
  19. };
    % u- s% @, m$ ?! F; }9 e
  20. ! ^2 a$ r9 C) c  y6 ]
  21. $worker->onMessage = 'on_message';
    1 K1 ]8 j; J' ^# G
  22. ) a# M8 g1 }% W! r/ J, Z( }
  23. function on_message($connection, $data)
    0 g- t( J( Z. r* N  Q. u
  24. {
    ' Y; Z# u3 A. Y" [# a4 O& C
  25.     $connection->send("hello\n");
    % |  m% L; [8 i+ ~& s: `8 ?
  26. }
    $ s* h7 u. \$ _- f! ]/ M: |

  27.   n, n- f! _3 o* G* A0 Y6 H7 |+ U
  28. // 运行worker4 m+ N* J4 p$ g1 J8 Q! \
  29. Worker::runAll();7 _  L1 j" M4 T; S* n9 V
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:% v& p0 e  q7 I1 |7 K

  31. : q0 n$ p- L4 Z/ O& j# \
  32. use Workerman\Worker;/ n" e( h+ Q- o6 C$ b9 m1 H
  33. require_once './Workerman/Autoloader.php';* X7 r7 M( a0 }. c4 Y6 B# n
  34. 4 Q, f( p; l& H9 P5 a4 q" E
  35. $worker = new Worker('text://0.0.0.0:2015');
    " g- x+ p- b3 u( R. g6 ^
  36. // 4个进程
    % s$ m( o% i/ x! ]$ V
  37. $worker->count = 4;/ s% w% B3 }$ t7 T8 v
  38. // 每个进程启动后在当前进程新增一个Worker监听
    3 ~/ s2 S  O7 i) p: o( q# P! j
  39. $worker->onWorkerStart = function($worker)3 q! L3 |* `7 p+ e  q2 C0 _
  40. {
    " L$ `' V6 \/ _: i# y0 v
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');3 F+ S$ i. l' @" B* f) A
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    0 z9 P& x( G2 \% m. Y
  43.     $inner_worker->reusePort = true;
    # Z$ N0 z$ l' L$ p; G
  44.     $inner_worker->onMessage = 'on_message';) E  g! t7 M: i$ {7 X
  45.     // 执行监听。正常监听不会报错
    5 W* n5 p1 H7 q, H( M8 |) U
  46.     $inner_worker->listen();; @$ z+ h& J; v  H( E0 y
  47. };- @9 M- ], P2 u1 F7 q

  48. 6 O; Q1 V' Q5 _
  49. $worker->onMessage = 'on_message';
    / M1 v" r7 a7 N: ^, F! J! y
  50. 9 e0 L. t) y' d" h3 [) e- o9 r
  51. function on_message($connection, $data)
    : ]% q, V& i+ p6 }" M$ `
  52. {( C' x7 P! R0 ?3 r3 H9 q9 X
  53.     $connection->send("hello\n");
    - }) G4 p; _1 G- [
  54. }  d1 t6 W. O7 \, i% O/ W

  55. , ?" N# b( B$ [5 h
  56. // 运行worker" M# d* V+ V: \& A7 f
  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
    3 i- {  ?- M6 u" _1 r" U  v
  2. use Workerman\Worker;- z" E0 \, ~/ g7 I( H
  3. require_once './Workerman/Autoloader.php';0 X3 R% O/ L( ]8 ?
  4. // 初始化一个worker容器,监听1234端口: d$ u3 g/ B* S! Z7 I
  5. $worker = new Worker('websocket://0.0.0.0:1234');" k; H5 }" R) Q. z, v, Y/ z

  6. : Y' `8 b3 M9 T7 Z6 M
  7. /*
    0 |. O8 G6 N  j. Q# h2 o/ H
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误7 |' }  x# u: W. w
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)- M1 m* Z0 }  r3 Y
  10. */. _8 D+ a, a' q9 m; C! H# u7 `" R
  11. $worker->count = 1;
    7 V) N/ F- C$ z0 }2 B
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口: F9 D( S! \; z  A/ x- M* `5 K
  13. $worker->onWorkerStart = function($worker)( m0 a) n2 B& x( M- h( K
  14. {/ L7 `+ P/ S, P
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    # H; l7 v( d; K6 j# G
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');% i; k( n5 Y3 C' h
  17.     $inner_text_worker->onMessage = function($connection, $buffer); D' s: x  A: s' c1 V
  18.     {% a) X; p2 U- ^- [! a9 u
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据7 A7 A- o9 b. B7 K8 A/ [' k
  20.         $data = json_decode($buffer, true);
    7 }% p( p' Z( m, C2 |1 V
  21.         $uid = $data['uid'];
    $ I  b" h$ Y# Z
  22.         // 通过workerman,向uid的页面推送数据
      t6 g/ b! M' B. B: k
  23.         $ret = sendMessageByUid($uid, $buffer);
    . D- G% Q) ]- x
  24.         // 返回推送结果) z9 r( P% S# A3 t; |* I
  25.         $connection->send($ret ? 'ok' : 'fail');
    * f% b3 R9 i8 y$ J, [
  26.     };
    9 y) j5 ^; ]% P4 s
  27.     // ## 执行监听 ##
    ' J! |% ]" o0 ~- r- i
  28.     $inner_text_worker->listen();
    8 H2 L; `9 t% N' h+ O
  29. };
    # W6 P& ?( N$ f, H& X7 f7 f
  30. // 新增加一个属性,用来保存uid到connection的映射
    * G0 U: l8 e' v7 W6 [
  31. $worker->uidConnections = array();7 j- O- w& {9 P2 o, s1 I
  32. // 当有客户端发来消息时执行的回调函数
    " v* s! z+ O7 V# ^, t# ]: k+ o1 P
  33. $worker->onMessage = function($connection, $data)
    : [$ ?8 t" E+ w' M  m5 w. q
  34. {
    4 g1 O) p7 h2 n3 f
  35.     global $worker;5 ]" Q) U. Z3 t$ l
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    2 m6 S9 Q, A$ z
  37.     if(!isset($connection->uid))3 \0 w* n7 T1 {% N( F
  38.     {" p& s5 e" D0 F/ A
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)+ J+ |+ B: s- P& h# G8 v
  40.        $connection->uid = $data;7 I* I8 @' Q" S- z3 m
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    # s( S7 O; n$ Y! `4 n: r' C. y2 j
  42.         * 实现针对特定uid推送数据
    - b. H* ]6 r% v
  43.         */" J9 Y9 t6 D2 L* [# U- E6 d  v$ t
  44.        $worker->uidConnections[$connection->uid] = $connection;
    - c0 e8 K1 p; l/ `" c2 a8 c: \
  45.        return;
    9 l/ `% g3 ~) ^8 o6 l
  46.     }9 Q) i* @& E. c; {; Z! j
  47. };& L* {* |5 n# H; J- x

  48. % l9 T" X: q0 Z! ~4 t" Q
  49. // 当有客户端连接断开时. N; z2 {5 |/ u4 n+ L- N
  50. $worker->onClose = function($connection)# Y6 Y5 }( G6 Q% \4 G. U. s
  51. {2 _7 }2 i9 N) E$ u
  52.     global $worker;
    # }  O" E2 [; q" d* a4 _; L# l( i# S
  53.     if(isset($connection->uid))4 `# L( ?7 p3 o* t+ H. q' w
  54.     {" `# [( H, q! T
  55.         // 连接断开时删除映射! L/ N/ D. s+ g8 p5 Q& m
  56.         unset($worker->uidConnections[$connection->uid]);
    8 a( o8 @) B% x, [4 W+ H* N" a
  57.     }8 q4 u! J2 L' ^4 T' g$ Q7 L8 c/ W
  58. };4 i2 Z% o6 U6 M( a$ E2 t* ?. v( O
  59. ( y# ~& A2 U4 B7 q# j4 {9 d' j
  60. // 向所有验证的用户推送数据* s! r% w1 ?# d, n/ r$ t
  61. function broadcast($message)
    2 v& }0 f9 N" a) i- m" ]# g+ O
  62. {
    ) @' D5 l3 c: }5 c1 Q: H
  63.    global $worker;: o% d% y* H" ?  r  }+ E$ o
  64.    foreach($worker->uidConnections as $connection)" ]6 F  p) m& I9 ^5 ~9 `
  65.    {
    + |& {9 c% f$ s! u; x
  66.         $connection->send($message);
    1 P/ z' X# K" A0 ?, z) \  l
  67.    }
    % |* e# z  }! O+ _$ z: N" Q8 X" @
  68. }
    + [4 y' ]1 B. \+ y) f1 b8 p6 J
  69. $ B" g. {1 @6 u# W. t
  70. // 针对uid推送数据
    6 A% P3 J) k  l. Y/ ^% H  q( I2 x
  71. function sendMessageByUid($uid, $message)  ?% [8 t1 N$ r2 a$ Q
  72. {4 V/ z) A: e4 m: a4 p
  73.     global $worker;
    . U. O2 X- v7 g$ Z4 {: c; i
  74.     if(isset($worker->uidConnections[$uid]))
    3 Z  O, m# F7 \% ^  M
  75.     {* s. |0 c( ^- X$ {, X
  76.         $connection = $worker->uidConnections[$uid];3 [5 b+ z$ [; l
  77.         $connection->send($message);/ b1 J4 x/ C! `: Z0 i
  78.         return true;2 ?' K6 ]; C0 X9 B7 _
  79.     }8 S- ?$ l, F+ Q
  80.     return false;! V4 q5 V. x$ t2 h- W; q
  81. }1 U) C1 h0 }. c1 Y
  82. 1 z2 g$ r1 e7 j) H. h' O
  83. // 运行所有的worker9 a- i2 N; ?: |. `
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');- h6 I' }5 O# w. X  d( p
  2. ws.onopen = function(){* K9 x+ @+ Z$ C; K4 [+ c
  3.     var uid = 'uid1';- e% I2 ], A9 h2 b4 Z( I5 H
  4.     ws.send(uid);
    - I* j1 \" d. n, x
  5. };
    . R+ B2 s4 y  u  c
  6. ws.onmessage = function(e){
    " S: S' Q: n+ q3 e9 `4 X
  7.     alert(e.data);
    & }8 r) E6 z: s) H  F! M* r2 K- \
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    * m% D6 n! ~# C1 \# M' Q7 g: a! e
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);: @0 N- W9 c) o5 W' x* f3 J  [
  3. // 推送的数据,包含uid字段,表示是给这个uid推送
    * ^$ P7 m( z% J; L1 W
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');
    / H7 r! O5 O9 Q( R9 B4 n
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    " `# I8 l# m" u, [2 W
  6. fwrite($client, json_encode($data)."\n");# x7 M, V: m5 ^' A4 s5 x. G, b$ I( I
  7. // 读取推送结果0 u" Z5 G1 \1 t$ b  {! K
  8. echo fread($client, 8192);
复制代码
+ Z% i8 Y( c# B) p# P' S- }
- g) k3 k3 f( R) {6 Z6 g
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 23:59 , Processed in 0.060134 second(s), 19 queries .

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