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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 16162|回复: 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;
    ! S$ v7 K* ?, Q7 v! N0 l
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    : [$ N: P& o; N1 H

  3. ; K. f: r7 t& u* g0 F
  4. $worker = new Worker();
    8 Q9 J  h- k9 F5 U
  5. // 4个进程# I$ \1 n* s. {  {- r
  6. $worker->count = 4;
    6 x& A/ N. E6 N/ R! P2 R
  7. // 每个进程启动后在当前进程新增一个Worker监听& W7 Y7 i8 a& k; p/ {  d
  8. $worker->onWorkerStart = function($worker)* ]: g' a2 W1 ]8 I* b" e4 Y
  9. {
    4 i* y* r9 K8 m6 i6 r
  10.     /**( S3 j  f1 m, d& t' G$ N( W: Q
  11.      * 4个进程启动的时候都创建2016端口的Worker
    - w* E- G+ P3 _/ B
  12.      * 当执行到worker->listen()时会报Address already in use错误
    ) K; m5 b' c4 m3 y; L$ \* o, p1 F
  13.      * 如果worker->count=1则不会报错: n6 o2 V+ ?& r' x3 y4 {
  14.      */, y( r3 ]; h6 c
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');
    + l1 i2 n8 z' l6 I& ]; \% E
  16.     $inner_worker->onMessage = 'on_message';
    9 I. f( z' }3 C2 a( @/ X& w; W! F
  17.     // 执行监听。这里会报Address already in use错误
    7 m. b$ B+ k& u1 V
  18.     $inner_worker->listen();8 Q# {# d4 [" E/ X& F
  19. };  Z% W) d7 B, ^" k
  20. & e8 K7 U1 E1 s! l6 o9 G" o
  21. $worker->onMessage = 'on_message';
      R6 d1 F( s: Q5 f$ _' c1 [: b% s" m
  22. . X5 D' X2 m3 _8 r3 S8 i* [) K
  23. function on_message($connection, $data)
    6 i. W7 s3 O. ~5 H
  24. {
      m1 p) K4 K- }5 I+ ^
  25.     $connection->send("hello\n");5 O4 @" A* |9 L+ O. A( y2 o3 y
  26. }2 `3 M  a" m/ n" x" k
  27. 9 t& d0 H% Z/ @0 @  Z
  28. // 运行worker: k' D* Y( e' N3 ^/ n, x9 Y
  29. Worker::runAll();1 b! e* c+ ]$ D0 o
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:' b7 I+ |) @% a# ?
  31. ( o2 }( C) g0 W+ v  p& v' u
  32. use Workerman\Worker;! L8 P+ ~2 H, O: x2 K& a
  33. require_once './Workerman/Autoloader.php';: V3 W% K5 |6 e. d% W, n$ ^

  34.   o3 _0 s  `: \+ Z/ h2 l& _
  35. $worker = new Worker('text://0.0.0.0:2015');
    4 t1 _! b7 E7 z# C  K$ X
  36. // 4个进程% ?' D4 ?* R$ y9 I9 Z$ z5 [" z6 w* ?
  37. $worker->count = 4;
    4 ^0 u' N+ D. _  R# j
  38. // 每个进程启动后在当前进程新增一个Worker监听0 V: w, x" f4 Y
  39. $worker->onWorkerStart = function($worker)
    & b; Z3 c1 z  }  L6 z, a4 u
  40. {
    9 w4 z- v5 F& @+ W# R# m
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    * C; u9 p2 G- p" K
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)% ~9 y+ A3 H+ J5 P* D* c
  43.     $inner_worker->reusePort = true;
    7 f& a  x6 a2 c9 b; g) [" f
  44.     $inner_worker->onMessage = 'on_message';
    0 P+ R9 L* W: F
  45.     // 执行监听。正常监听不会报错* G  V( R1 R5 X- d6 M' Y) y; J$ r
  46.     $inner_worker->listen();/ c: s; {* v% ~- D0 c: e
  47. };7 u/ y" w6 Z- b6 E0 ?: ]
  48. % Q$ f' T% W: B0 i* P* z) R/ s
  49. $worker->onMessage = 'on_message';
    7 j/ ?# W# o; t
  50. - R6 d5 \4 j, W
  51. function on_message($connection, $data), c% _) G* G2 b- q, h
  52. {& U% X% ]$ }( a  @1 O
  53.     $connection->send("hello\n");: y8 k& u7 `) B  g2 J2 ~4 K$ D
  54. }' Q1 I0 X! a0 z2 G, j( ~
  55. : I3 y# a- {# b) `5 l
  56. // 运行worker6 I9 ^  c* f. x0 P
  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/ |6 u& A: g- E+ U
  2. use Workerman\Worker;
    / g% i: q; R9 S: o& b; C- K
  3. require_once './Workerman/Autoloader.php';
    - J8 r, n/ n' j. w/ s$ Z
  4. // 初始化一个worker容器,监听1234端口/ B" I' o5 Y$ k- p1 J
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    3 I; l" M  f# r9 A3 Q2 U$ [' r( N/ a

  6.   N5 k. x  l! n/ O+ ^) i
  7. /*
    7 T7 \5 G9 J+ }% I1 R+ m
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    0 Y2 o4 \% D; w. [5 Q# i( z) r1 d3 b
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)( `. h6 a0 z  J( Q. _9 Z& Y' ]
  10. */
    ' v- b1 ]3 d5 m& _" w% A
  11. $worker->count = 1;( t- {' d2 g, S5 o3 u
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    - q! ?; s7 \4 a$ A+ |
  13. $worker->onWorkerStart = function($worker)) P/ M: X# x( Z
  14. {
    3 r# W  C& v! t1 J* y
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    . A2 f6 z) N8 f6 w0 j
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');
    3 K" C, m. P  j
  17.     $inner_text_worker->onMessage = function($connection, $buffer)
    2 _3 u& T7 L; b9 [3 z
  18.     {+ o1 r) o$ A, L
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据
    / ^% X# n7 J1 Y6 _
  20.         $data = json_decode($buffer, true);+ X0 B7 h: k. r: b
  21.         $uid = $data['uid'];
    2 p; B7 a4 N& o, h
  22.         // 通过workerman,向uid的页面推送数据8 w9 E+ m2 h# X6 e
  23.         $ret = sendMessageByUid($uid, $buffer);8 |) p! h+ r5 D
  24.         // 返回推送结果
    & Q7 Y; ~1 L9 m' d
  25.         $connection->send($ret ? 'ok' : 'fail');4 Z; x. y& j6 u: Z$ Q  p
  26.     };
    ) J$ \4 a, }3 v0 ^7 y
  27.     // ## 执行监听 ##' y9 A3 J/ A6 s$ p& \- G" P: p
  28.     $inner_text_worker->listen();0 ?  ]) |7 N. h* r7 ]* b$ g0 P- I5 T
  29. };, [3 o+ H: K! n& A2 m2 u* n
  30. // 新增加一个属性,用来保存uid到connection的映射
    + h; {: a$ L7 t- S. `8 d: D
  31. $worker->uidConnections = array();
    5 X- e8 D. |* C! m+ p/ g
  32. // 当有客户端发来消息时执行的回调函数
    5 C  s) X2 m7 q
  33. $worker->onMessage = function($connection, $data)
    + h4 l2 k/ H3 K& w( A
  34. {' S& R  J1 J: V
  35.     global $worker;% {4 ~$ r. k/ p( g/ N: ?
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    9 F1 G8 @" V  C! \
  37.     if(!isset($connection->uid))
    6 \; }" v1 e$ F0 j$ |* t/ y9 P7 T
  38.     {
    $ ?& U" L. v4 V( A6 E
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
      n+ y+ E- H! I% `8 S
  40.        $connection->uid = $data;- `: B* O  g  R& U3 y
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    5 V* W4 V* ~6 c4 i9 c( _* l
  42.         * 实现针对特定uid推送数据7 w% }" C+ Y1 s+ \: U
  43.         */
    + N5 U, e7 u$ B. x8 ^2 n
  44.        $worker->uidConnections[$connection->uid] = $connection;
    8 Q4 Z- F6 Q8 e+ X
  45.        return;2 o& j; l$ E! H$ B- `1 p
  46.     }
    4 {$ h. V2 g6 n* {
  47. };
    4 o2 i) |- U; A/ @" K
  48. 1 h: h0 N" _& T: n! k
  49. // 当有客户端连接断开时& z( H8 K! d! Q9 I3 L1 d& }# W
  50. $worker->onClose = function($connection)& M3 L/ I' O) Z5 |1 H4 g3 }) i' j
  51. {& E5 f! s. k3 [( |8 s
  52.     global $worker;4 }* k9 @2 V# T
  53.     if(isset($connection->uid))
    . E, ?8 H- h0 m9 n7 h$ {
  54.     {1 M+ y. u4 \" y7 y! S
  55.         // 连接断开时删除映射, c3 ?- y1 Z" a2 B" l6 L
  56.         unset($worker->uidConnections[$connection->uid]);
    * F8 `7 y9 m  g* _. C- Z
  57.     }6 ]) Z! s1 X; }% Z9 f7 c3 G
  58. };
    " _' k* l5 M/ s$ }8 ~& {
  59. . n* L2 X5 b/ |5 G. U' m
  60. // 向所有验证的用户推送数据. C9 I: S6 C/ E. d" ~
  61. function broadcast($message)& X. v6 O% S1 J9 H2 m% N; d1 L6 N' K
  62. {
    ! k2 H' k  @# J* t2 @6 h; m+ U. z
  63.    global $worker;& n7 o) p( a8 A( H1 v5 D
  64.    foreach($worker->uidConnections as $connection)) ?; q2 ^& b& E& Z
  65.    {4 h. ~6 Q( A/ b6 N4 E% D+ K3 {8 x
  66.         $connection->send($message);
    0 y, z* u: X0 [4 }( s  n
  67.    }* L6 {2 Y: a* L
  68. }
    + L8 q( U% U; ^% Q

  69. % ?9 a& w9 a3 X$ B
  70. // 针对uid推送数据
    ' Y; G3 y8 A) ^) h
  71. function sendMessageByUid($uid, $message)" `2 x4 }; k; x( F9 O% E- O8 D
  72. {
    ) [' i% m+ P# n/ s/ L* h
  73.     global $worker;5 f$ Q9 D! n6 G, \$ x' G
  74.     if(isset($worker->uidConnections[$uid])); q. Y2 i5 Y7 N: k
  75.     {
    0 a8 ]8 s/ {4 T, q
  76.         $connection = $worker->uidConnections[$uid];
    ) M% O3 O9 A1 f% N  d2 l9 G
  77.         $connection->send($message);/ n4 K7 l1 Z' {6 @; H  W
  78.         return true;
    % i8 r( T& f' n* Z, M& W+ L
  79.     }
    # d1 n& V' ?& N0 J
  80.     return false;
    , e/ u0 R! b( e/ r3 P
  81. }
    0 G+ R& m% n! x& X

  82. ( W0 @$ l+ n5 T# X  x
  83. // 运行所有的worker  b. C/ I+ T8 k4 K1 a; f: }
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');( x5 A" P& M8 k
  2. ws.onopen = function(){; q& Y( f4 a  C: W( f
  3.     var uid = 'uid1';  O+ H- O+ l# g: \# P
  4.     ws.send(uid);& C6 K0 R# T. Z% v1 V& a
  5. };1 \+ Y% `  ?7 V0 J/ |9 A
  6. ws.onmessage = function(e){
    1 Z# a! x; u3 ~/ u7 ]
  7.     alert(e.data);
    ' w! x; f, o: o$ G1 R# ~
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    + z2 z  A1 ?, g6 x. M5 E
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
      e) X$ l4 O- Z' _: m
  3. // 推送的数据,包含uid字段,表示是给这个uid推送
    6 \" y5 x! X5 P. O! L
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');* n! J. }1 O, s5 x5 p% b: p) f
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    9 A  |9 D% O+ f/ N& c3 D& \4 m
  6. fwrite($client, json_encode($data)."\n");! p/ u) U# o; K3 _% s: X
  7. // 读取推送结果0 V+ m' }) S8 [/ F' H. n& B
  8. echo fread($client, 8192);
复制代码
; }3 h, _( V0 Q

! Y3 W. _8 ^2 q# I) M( h
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-9-21 05:17 , Processed in 0.046396 second(s), 19 queries .

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