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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 16193|回复: 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;/ K0 e4 V9 y3 g! [2 w0 A
  2. require_once __DIR__ . '/Workerman/Autoloader.php';& y! Z+ {( |1 `- {

  3. ' X; i. |. W+ z; V3 e* `
  4. $worker = new Worker();
    9 V  C+ Q- r7 A+ o4 P8 L4 W" I
  5. // 4个进程+ `' U1 P# d# ~. z
  6. $worker->count = 4;6 p' |. R& w" @
  7. // 每个进程启动后在当前进程新增一个Worker监听
    - {! _7 u$ i9 ?2 u
  8. $worker->onWorkerStart = function($worker)
    / X2 B! `4 e. A5 ?$ U4 P3 _- I" g" z4 ]
  9. {6 d2 z; _/ Y! F, F) s, N; L. p0 ~
  10.     /**( h+ Y- m$ ]* _+ R& D, E9 j
  11.      * 4个进程启动的时候都创建2016端口的Worker4 ^! i; ^( N. `" d4 T
  12.      * 当执行到worker->listen()时会报Address already in use错误
    ! N8 d( _# N% V/ c( b7 x, i
  13.      * 如果worker->count=1则不会报错
    + Z+ {5 }: W& a. q- }
  14.      */$ x! `$ G0 w# `% J# m( S( }
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');& S! b% Z2 z. c0 Q; K  _
  16.     $inner_worker->onMessage = 'on_message';) Z1 C; |1 ^5 P: u6 K! Q; d9 P, ]
  17.     // 执行监听。这里会报Address already in use错误
    ) O+ A* U5 d5 a
  18.     $inner_worker->listen();! L4 M! ]. o' R1 r2 A2 |# _
  19. };
    & k2 _0 I- m4 l8 n
  20. $ u" r, C1 b) m% _
  21. $worker->onMessage = 'on_message';/ G0 R$ M1 Q( @* j0 ^4 _( P
  22. " m. r+ g4 z& I& c5 w$ U, I/ O8 P
  23. function on_message($connection, $data)
    ) B9 ?$ I+ G; ^/ D7 p2 w
  24. {6 p& D9 L/ A5 Q3 i* k/ a5 r+ c( R
  25.     $connection->send("hello\n");% S1 G. n8 r2 N4 S* l1 [
  26. }
    5 d; I' i$ z! o
  27. 0 I0 i' g: ~7 y: y
  28. // 运行worker
    & U. z. {! c- B1 `: {2 G  G
  29. Worker::runAll();
    2 O! |0 x% z4 U& j+ |% y% p  C
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
    " \* o4 @, ?) O& Q: @6 b7 u

  31. 3 }' c* n6 s0 k1 ^
  32. use Workerman\Worker;
    : \! o& ~+ Q$ r" N# A% o7 z
  33. require_once './Workerman/Autoloader.php';* x& j/ W+ g$ j4 {3 j& X: H

  34. * ^4 x; ?* ^, W; {% B
  35. $worker = new Worker('text://0.0.0.0:2015');
    7 e- C9 d0 H9 ~1 h: z
  36. // 4个进程
    ! |; h3 Z: z0 o. z0 P  n
  37. $worker->count = 4;
    % E/ Z$ j- s% u; e
  38. // 每个进程启动后在当前进程新增一个Worker监听
    $ p+ X. L. ~. A8 [; A$ P1 V  W
  39. $worker->onWorkerStart = function($worker)
    * s  h+ Y) K+ x6 h8 p& }) }8 a
  40. {# ~" f/ E/ h2 c! p! A$ j6 c
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');2 l) H6 F! g* }) Z- q5 V
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    & J: P6 |0 r& W) Y; T/ y# P; ?
  43.     $inner_worker->reusePort = true;
    / l% k2 |) `3 ^
  44.     $inner_worker->onMessage = 'on_message';( t) G$ b, J( B7 D
  45.     // 执行监听。正常监听不会报错% ~8 u5 Q. T. O' @) R$ Z
  46.     $inner_worker->listen();
    , }5 }! E: q, B/ R# }3 C
  47. };6 l. w* V" }+ C6 v8 ?
  48. 7 y! B5 {( Y( P8 W* w
  49. $worker->onMessage = 'on_message';
    4 X( H1 R7 x* k9 V- {" p. O
  50. ( l7 n) @" R/ G+ x" d1 C
  51. function on_message($connection, $data)
    9 @: J6 R3 }  d( m
  52. {
    3 {1 T+ S; l" m. k% ]3 z; s  Q
  53.     $connection->send("hello\n");- z  H& \8 H; U9 ~" m! Z9 B( |
  54. }3 A' A0 @; B( L) y, l# Q
  55. & D- }" n* ~; }& `
  56. // 运行worker4 a; Y2 Y! T/ W. P! t' H! k# @6 E( D
  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: c- R( Q# P9 {) s. |1 T# a
  2. use Workerman\Worker;
    + A2 \5 a8 A' J% e- \/ F! w5 A0 j
  3. require_once './Workerman/Autoloader.php';
    6 R/ ?4 x; j9 S$ f6 L
  4. // 初始化一个worker容器,监听1234端口' n7 v& k% ]4 J
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    ! c/ a0 x6 l8 @
  6. 0 z+ D% O9 ^; l8 p
  7. /*: z/ `& |# |6 g6 W
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    8 y; g( X( g4 x7 I1 A
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
    / H; {9 [1 [( x/ G* x+ R0 w
  10. */9 ~% I) v4 |* f
  11. $worker->count = 1;
    ) ^3 Y2 X! O* v5 p9 {& @
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    : j! C! Y2 }$ y7 O( k
  13. $worker->onWorkerStart = function($worker)
    $ s; t+ v6 c: m! [
  14. {
    + w- Z( C8 s0 j& b( i5 X: h$ N9 ]
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    8 ]+ M6 t% ?$ i; g
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');
    # F4 i: \3 x, h/ r) o
  17.     $inner_text_worker->onMessage = function($connection, $buffer)0 w  V5 D, K: F4 b5 {, A
  18.     {
    % ]/ o2 N0 h2 u' k2 r
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据
    0 E/ e+ }% b( @9 n% N  S. y$ ?
  20.         $data = json_decode($buffer, true);
    " h) h3 N7 }5 e/ x& h, g
  21.         $uid = $data['uid'];$ D* K9 ]8 s! w4 `# X$ v7 F4 B
  22.         // 通过workerman,向uid的页面推送数据" p+ R( ?) [4 M
  23.         $ret = sendMessageByUid($uid, $buffer);
    ! V2 L! g/ X' T* }6 H/ S1 \0 s
  24.         // 返回推送结果' D1 x1 `6 L! x8 o  T, E
  25.         $connection->send($ret ? 'ok' : 'fail');
    6 a8 w  ~" h/ F  N. P1 z
  26.     };
    - Q0 z" j7 Z9 _9 G
  27.     // ## 执行监听 ##
    $ J+ O7 D. V0 R8 L% h: c/ A
  28.     $inner_text_worker->listen();
      Q8 E. s6 Y1 k- K8 m* W! F: O7 B0 ~  {
  29. };, X7 a8 L  n4 R, d8 c
  30. // 新增加一个属性,用来保存uid到connection的映射: E8 r* z" u$ D
  31. $worker->uidConnections = array();
    " Q+ A$ h9 p9 t. g8 P5 F( n$ E
  32. // 当有客户端发来消息时执行的回调函数
    8 a; o, w" _5 n+ r0 ^: ?4 R
  33. $worker->onMessage = function($connection, $data)
    / s2 N2 L5 W0 l" V( i+ b1 V
  34. {) T  h7 l! n' u; b1 B7 ]$ g+ `
  35.     global $worker;
    + Q# t% Z6 j- b' V6 C
  36.     // 判断当前客户端是否已经验证,既是否设置了uid! ~9 }9 g  M! @* X
  37.     if(!isset($connection->uid))
    , S- k9 u" R5 t1 t7 c5 C( R
  38.     {
    ! u) O5 b  O* Y  G
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证); i) |. I$ p9 B7 p* P$ i+ `
  40.        $connection->uid = $data;
    / q, P6 V( q8 M0 x# ]" t  f( ^
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    2 _9 m, L3 y  W- z0 u5 ^6 }% u
  42.         * 实现针对特定uid推送数据
    0 z- u' T3 p4 \8 J' a( m
  43.         */
    + w& }) M# o; D& e
  44.        $worker->uidConnections[$connection->uid] = $connection;7 P/ y0 g2 z& F+ v* |. Y
  45.        return;$ `7 i7 G# u( }( L; y9 K7 q
  46.     }
    , A" y: a* k2 m. ~, x( D
  47. };
    : t9 }$ z0 v1 N( u. v8 g
  48. 2 ^( f, R9 f* J: Y4 q
  49. // 当有客户端连接断开时8 @; \8 L* x) p" T! j: A  O
  50. $worker->onClose = function($connection)
    ' O8 T2 h  o+ b. C: d8 x
  51. {
    3 A; p* o( L4 i1 {& X
  52.     global $worker;
    9 i& S5 S! Q' Q- P2 e/ @1 y
  53.     if(isset($connection->uid))
    7 k4 ^/ I8 O& R4 e
  54.     {
    . q- t" K5 d) `) a5 q
  55.         // 连接断开时删除映射
    ) ?+ A+ ^8 _% {: M' n
  56.         unset($worker->uidConnections[$connection->uid]);
    ) l, y7 Z6 \# T; V/ }
  57.     }- D$ W' A, _3 n4 {
  58. };3 e. y" }% H# J  [: Y' S* ]- b

  59. 1 |7 S" ~  I( @7 g6 y
  60. // 向所有验证的用户推送数据
    : O5 j2 S+ h9 U- w; M* `& y
  61. function broadcast($message)
    ) t# U; I; @2 f5 t: ]6 E
  62. {! v# \$ x, P' o! y& N0 Z9 S6 X
  63.    global $worker;
    # Z" ]) k  S: M( K, q# L
  64.    foreach($worker->uidConnections as $connection)
      o% X) ?' a7 W# v& t1 b# u
  65.    {5 D: F6 }. x8 O7 k0 O; u
  66.         $connection->send($message);* r8 O; R; _% p6 `- b- L& o$ c4 x
  67.    }- c$ Y3 h8 {8 o2 @% d- j: E( d
  68. }
    & x) f' r+ g$ S2 C, T- i3 ^* s% l

  69. ; [2 b# F+ g# b! A
  70. // 针对uid推送数据
    9 z1 U2 g: s, m8 h9 Q
  71. function sendMessageByUid($uid, $message)8 d; c# \+ K9 {" `
  72. {0 A" u/ b1 V  B$ s* b' {
  73.     global $worker;
    $ K& q3 I* T1 y- m% Y
  74.     if(isset($worker->uidConnections[$uid]))
    + \" U3 a$ O( m+ [! s4 C$ f3 b* i
  75.     {
    1 g# n. z0 D: z6 x& i2 }
  76.         $connection = $worker->uidConnections[$uid];$ W$ z" L, D  \2 c
  77.         $connection->send($message);# O2 q0 {; }; {7 b  i: `0 h/ E
  78.         return true;$ K8 c4 S  e  Q8 A
  79.     }. O7 L% ~* k" }, T1 E" V- Z# k- I/ Y
  80.     return false;
    . K7 B# H8 s) p/ E$ k6 Z
  81. }  B) N" s+ ^4 O% Z- i+ h

  82. . T! z6 F2 l! e0 ^6 E5 _1 I  Z
  83. // 运行所有的worker: Z4 H8 q3 f/ V6 |
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');
    % @# g. h/ z: j( Y3 j+ K
  2. ws.onopen = function(){
    2 _) i  Q3 A' z& w- Q2 y
  3.     var uid = 'uid1';
    8 p  \' c3 s5 J9 V& j
  4.     ws.send(uid);
    * Y; v+ l4 r- \' V7 Y3 I
  5. };6 x  X& Q0 i' b& D) H; E( R
  6. ws.onmessage = function(e){+ I+ ?/ H+ w, e: B8 n7 i
  7.     alert(e.data);3 d: B  W; C- N# r  r2 w! @0 X8 j/ D
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口, X% J2 D5 g% _- o+ x- L$ l
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
    6 _6 D/ q/ ?# S7 _
  3. // 推送的数据,包含uid字段,表示是给这个uid推送
    ) ]% T$ q3 }2 F7 T* V; ^
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');* ~7 N3 e0 X7 A9 h" _; R% }
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    ! [" l! r5 L& |8 t6 Q: {
  6. fwrite($client, json_encode($data)."\n");
    0 a! t/ X" s8 Z/ r/ Q: A3 v" }
  7. // 读取推送结果
    % @; N- Z9 R: |. W
  8. echo fread($client, 8192);
复制代码

# |) y6 m. o, G. F9 ?. ?4 O7 J5 ]! C) H
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-9-22 07:00 , Processed in 0.070106 second(s), 19 queries .

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