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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15612|回复: 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;
    ( `* N$ w/ t2 J4 K
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    ' v( U+ \, ]) \* u' B2 ~
  3.   {6 t: [8 e7 }$ s+ G
  4. $worker = new Worker();
    $ C' M) i( H* v# f$ H6 |8 E
  5. // 4个进程
    ' K$ \  m+ ]0 C4 i5 Q
  6. $worker->count = 4;
    3 M) Z$ @: C0 b
  7. // 每个进程启动后在当前进程新增一个Worker监听
    2 o! h* s+ q0 D! E+ X
  8. $worker->onWorkerStart = function($worker)+ n) q" Y7 k% F# ]
  9. {
    & n6 n7 Z1 r* ^( o$ R& ]  a9 y6 b
  10.     /**
    . n" s6 E" n0 R9 h/ b! W6 s6 x% m
  11.      * 4个进程启动的时候都创建2016端口的Worker3 c# d- H' |4 Y3 t! A3 a% K
  12.      * 当执行到worker->listen()时会报Address already in use错误
      q7 i% v/ H3 G/ I  K/ H. o
  13.      * 如果worker->count=1则不会报错
    5 a* Y3 d  Z& m' e. m& H: a& A
  14.      */
    / @2 W0 W9 O4 k, u! @8 X
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');( h% M. S% V1 Q+ |
  16.     $inner_worker->onMessage = 'on_message';
    ' h9 _& f5 d0 j0 W6 h1 o
  17.     // 执行监听。这里会报Address already in use错误
      S7 W" e0 \& N( X
  18.     $inner_worker->listen();' b; H! ^% X3 X" z7 \
  19. };4 U# A( }. |- b! A# k
  20. 0 l' w" }( C& L0 {9 c; J. f0 i! L7 z
  21. $worker->onMessage = 'on_message';+ c* c4 ~: B, C0 q: i4 d# A" b. L

  22. 5 L( k& Z. i. n$ Y. }
  23. function on_message($connection, $data)
    , |* e6 Q. O  {' U9 c
  24. {9 U# J6 F7 ~* w
  25.     $connection->send("hello\n");
    " A2 K! i) a( z0 c8 T
  26. }
    ; P! e( S* H8 N; Z& Q* Q( ?# T# w7 E

  27. ; C% [4 I2 T: b' H+ h1 N0 f8 _9 Z
  28. // 运行worker/ Q- k5 _3 g; Q: M5 w: D
  29. Worker::runAll();
    ) F+ O& B1 y$ Z5 Z0 S1 f
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:0 H  l' j% Z# n. f9 ~$ e
  31. - n  j9 q8 i7 P4 t
  32. use Workerman\Worker;
    ( F% k' `* @/ N
  33. require_once './Workerman/Autoloader.php';
    / p1 u7 W' i/ v8 c! p
  34. # J6 W* T6 X% F
  35. $worker = new Worker('text://0.0.0.0:2015');. [; ]# [8 }- `) F+ t8 R
  36. // 4个进程
    2 h7 n  N& E6 V$ z& D
  37. $worker->count = 4;
    6 w% p: A+ u7 x2 E
  38. // 每个进程启动后在当前进程新增一个Worker监听, n& G/ A% z4 W/ k, J; U
  39. $worker->onWorkerStart = function($worker)
    * f# ?- u$ x# u& I8 u  I
  40. {
    ! I" T4 u, j9 L& h
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    # S2 g. ^5 t7 T9 p- t. q; y
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    : J. ]) e- Z/ y  O. }: }
  43.     $inner_worker->reusePort = true;, g: M( v. x4 u+ Z, Y' y/ s4 I
  44.     $inner_worker->onMessage = 'on_message';
    % J* A0 G. u( m) w* p- ^
  45.     // 执行监听。正常监听不会报错) |. ?0 K3 t. f. [7 Q
  46.     $inner_worker->listen();
    ; N/ v, k6 H% M
  47. };
    % u6 T, M! i. ~: d  A

  48. 7 g1 F7 v" m. z+ ?: t; n* A2 m
  49. $worker->onMessage = 'on_message';% `; m/ S: S! h% ~) D+ V

  50. % m. U9 y' f) c
  51. function on_message($connection, $data)8 N' `3 {( {: s6 W' Q2 |% D* h
  52. {4 X9 e2 x/ c( ?: E
  53.     $connection->send("hello\n");
    % f0 x4 B  r! n5 R
  54. }  \, Q" P& ^6 q, A4 o5 ~; d
  55. 2 s6 ?% v% _2 `6 @" ]' k; Q
  56. // 运行worker
    9 {5 E5 y" U* z( I$ W0 h
  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
    8 |& {( j3 ]5 Y# |* J
  2. use Workerman\Worker;
    $ u( ]0 @4 [! M4 s0 p
  3. require_once './Workerman/Autoloader.php';
    " J: S4 O8 M( q4 {' H, h/ i" t( b1 e
  4. // 初始化一个worker容器,监听1234端口
    - y9 A3 p0 q' r5 E
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    7 v& D, t; c9 p" f% h& ~% ]; Y
  6. : x/ A1 Z/ Q5 w! x
  7. /*4 v  v- ?$ I2 ?! H: B  b
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    0 f: l3 \5 p6 l$ G& X5 {
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)  ~% g+ l8 }1 ]0 f; w8 _: d8 r
  10. */* H6 M, u( {* u3 ?; Q
  11. $worker->count = 1;6 j( F, u) Y  l5 G* w$ }  c
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口5 T# v8 w1 F" f2 C9 C9 L
  13. $worker->onWorkerStart = function($worker)
    : a) F) L0 x# B+ K
  14. {
    + g2 ~2 X/ _( E+ |. v
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符1 x2 S, l8 X: g6 W: W. Z" C$ G
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');" s* E2 |& G! k/ Z8 e0 ^" f9 @
  17.     $inner_text_worker->onMessage = function($connection, $buffer)
    5 V2 m1 e" j7 f! R. C: V
  18.     {
    3 Z. T6 P. u! k  W, m& o
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据% A9 D0 O1 a; e
  20.         $data = json_decode($buffer, true);
    5 R3 ]6 r# \  N7 b; u
  21.         $uid = $data['uid'];
    9 [% F4 f9 Y  q' v: c
  22.         // 通过workerman,向uid的页面推送数据
    , n& I! Z  ~" t; D; E
  23.         $ret = sendMessageByUid($uid, $buffer);% r: A# F" a0 D& p! `
  24.         // 返回推送结果
    ' ^# j( e9 p( ~# G4 G
  25.         $connection->send($ret ? 'ok' : 'fail');
    * l7 u: {6 Q9 S0 P* b' C
  26.     };' b( G$ y% B. @% h, e
  27.     // ## 执行监听 ##, Z! \$ p, M( b5 ?! s- H
  28.     $inner_text_worker->listen();  e3 b/ C$ L  U; g$ N  s' n
  29. };! V9 u- c/ d) Y  W
  30. // 新增加一个属性,用来保存uid到connection的映射
    / g8 B) k+ \2 R! e
  31. $worker->uidConnections = array();  }9 c: f5 j1 L! x! C: c
  32. // 当有客户端发来消息时执行的回调函数
    ' q: W4 Z6 i# C+ g7 v' i5 N5 z# [
  33. $worker->onMessage = function($connection, $data)% l% C% ^4 p5 A$ G
  34. {, D( e2 f" r- J- g
  35.     global $worker;/ Q7 R% q# {0 D: ?$ w
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    8 Y  H( |5 L3 z0 x- {. r7 O
  37.     if(!isset($connection->uid))
    , {- C' \" P  C& u9 n0 W% I
  38.     {4 g, _5 _" T7 z% R% S( D3 c
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)# [6 B9 T& U. h) c
  40.        $connection->uid = $data;
    2 ~7 Z9 ^8 u" O* [: [0 Z6 Z
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    ; u% O/ w0 k7 r
  42.         * 实现针对特定uid推送数据2 y9 Z. e. W$ E( K
  43.         */
    6 N3 O# G' I4 S5 n9 W9 h' Y0 U
  44.        $worker->uidConnections[$connection->uid] = $connection;
    + v0 v6 f- M4 H6 w0 Y
  45.        return;
    5 f% e, E7 I, j0 R- W$ ]
  46.     }. J# y9 \# B% w% }" h8 Q
  47. };/ o9 l' r* h. L1 m0 F5 P( O/ O

  48. : k; v8 y1 K' _! j
  49. // 当有客户端连接断开时) ^( m) K1 a6 f! }/ D
  50. $worker->onClose = function($connection)0 Z8 t" o+ w) [! I* j- O
  51. {
    8 K' P" }0 s4 b% K6 S, @
  52.     global $worker;7 O; ?9 t3 \! a# g/ G
  53.     if(isset($connection->uid))7 z6 w" [  B" X9 {3 [
  54.     {
    ) l2 o2 n& f9 @3 A, J/ \
  55.         // 连接断开时删除映射% h4 r0 u0 l+ q$ O2 m; C" _
  56.         unset($worker->uidConnections[$connection->uid]);6 h1 |9 A6 U9 r% i
  57.     }
    4 R  `" B; g6 c" v2 `! G
  58. };6 X  C) Z& n0 |) C2 Q; v. Z
  59. 3 V% y- C; K: V5 X4 w% d7 Y
  60. // 向所有验证的用户推送数据
    $ q6 X0 p* C3 a: x# r( m5 ?& r
  61. function broadcast($message)7 r7 d  y' |9 I7 d+ i, c
  62. {& a! V  x3 i1 l7 o; {; u
  63.    global $worker;
      v# S; i# E9 U1 R$ `
  64.    foreach($worker->uidConnections as $connection)
    0 Z% O' k* L! d1 L' W
  65.    {$ p+ ^! I2 o1 S$ S2 \
  66.         $connection->send($message);" f% F/ t; H" M
  67.    }( Q% q  }0 `+ c2 y/ }; `
  68. }' @, V( ]$ v# G* D

  69. % l% U1 J$ o3 g) b/ h
  70. // 针对uid推送数据
    3 e* o0 B9 O2 G; {, Y) S# {
  71. function sendMessageByUid($uid, $message)+ A( L& @0 M0 h
  72. {
    ; w" G9 [+ T# E3 n2 N' ^7 {
  73.     global $worker;& L% b9 ^: V! |) ]2 }% h
  74.     if(isset($worker->uidConnections[$uid]))( ^3 Q8 h, v* }, C
  75.     {
    3 l# E& f% z' c. S* g! i3 B
  76.         $connection = $worker->uidConnections[$uid];
    " v" Q/ |9 {) d; G- N# [
  77.         $connection->send($message);6 z( v  F, U& O' P
  78.         return true;
    5 b9 g% K  @: f+ i1 j; d+ A
  79.     }+ N' u' {" E0 ~5 J2 k6 C* o
  80.     return false;
    5 r1 R7 E+ w. j1 b" \
  81. }1 B1 P  o0 i( P. ?4 k% }

  82. ! @; Y0 h, n- _% t; T( D
  83. // 运行所有的worker
    1 _2 `5 P) h" A+ Q! H" |
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');- M5 B  L7 u, e1 b+ t! J5 t
  2. ws.onopen = function(){
    + a4 N( S' {; U1 ]
  3.     var uid = 'uid1';! _1 [# }# W. Y# y. [+ R/ a
  4.     ws.send(uid);
      R/ i- @. Y( U5 C4 ?
  5. };
    8 w( l& t$ R4 y  p0 v8 c' A
  6. ws.onmessage = function(e){
    4 R+ ~! a% f* u6 }. f/ K  o
  7.     alert(e.data);
    6 _4 X# r5 T, X% E' }
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    + N( [( j8 K8 q
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);* w6 l4 h6 ]9 ]9 Z& V$ j5 i# i" w
  3. // 推送的数据,包含uid字段,表示是给这个uid推送) N& M+ y0 K5 T" [9 L) Y/ }
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');
    ) Q9 P! m$ c; d- b
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    ' a" {0 h. n. p2 g& f: c
  6. fwrite($client, json_encode($data)."\n");
    $ f; Z+ [  H+ a+ p3 a
  7. // 读取推送结果
    - j, z2 A9 o5 A& J, p+ ~4 S  @# t
  8. echo fread($client, 8192);
复制代码

. a0 a2 ~9 X. N2 ?* P6 S# B" }. @
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 10:06 , Processed in 0.060002 second(s), 19 queries .

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