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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15613|回复: 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;0 N. g1 E+ Z. s
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    % V) V+ @& _: r& C3 u3 j

  3. 5 O( r# u5 [) v: j1 S& w& ^
  4. $worker = new Worker();
    : u, a# o- U3 e" m. ]5 U) \
  5. // 4个进程
    2 J& \# C; I) }( E$ C
  6. $worker->count = 4;/ t* I+ g7 M; v/ L
  7. // 每个进程启动后在当前进程新增一个Worker监听
    3 |: p' ^7 o$ M
  8. $worker->onWorkerStart = function($worker)
    % G' Y+ Z7 v& e/ e
  9. {
      {* E0 r0 z/ v+ ~
  10.     /**% J/ \9 K: `2 i
  11.      * 4个进程启动的时候都创建2016端口的Worker
    3 {! r, f7 K; @/ A. C
  12.      * 当执行到worker->listen()时会报Address already in use错误- p8 `" p. x7 U- f! ~" E8 _& K  [+ V
  13.      * 如果worker->count=1则不会报错0 l/ Z- f  m. _$ A' X
  14.      */$ X4 }) M9 l, z, f& B
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');
    . e) r! @8 o% h' N7 D, G4 y
  16.     $inner_worker->onMessage = 'on_message';
    4 q1 H6 I5 ]/ f1 |- R0 ?% j
  17.     // 执行监听。这里会报Address already in use错误; A9 K/ a5 r# ~- M1 a- w
  18.     $inner_worker->listen();) C4 I8 E# d* |8 _& g
  19. };
    2 Q# W8 l" p! e  H% n& ~
  20. / U' k- r4 |) [: {# R5 w
  21. $worker->onMessage = 'on_message';  x2 w& Y& ?3 y% r! ]0 ?

  22. & a: X1 F! q3 j; ~
  23. function on_message($connection, $data)3 v3 t" g/ ~. s6 @2 T+ o0 e9 K2 Y. M
  24. {
    0 \% N' j1 f' z
  25.     $connection->send("hello\n");
    & G4 v. r3 ^. I
  26. }) T4 {% S/ u5 e
  27. % X1 N0 N1 b/ ^9 n- O
  28. // 运行worker# J- ^/ l# F2 k9 X
  29. Worker::runAll();
    % f/ w# V7 ]7 g' ~' O
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:' L7 W; m+ Q5 J- M

  31. & @6 _4 D: T$ c% F  ]+ O1 ?3 P
  32. use Workerman\Worker;
    . S4 p. D) Q# H7 [1 l( P
  33. require_once './Workerman/Autoloader.php';; k% n+ Q3 Q6 ^6 [

  34. 6 F( ]8 J7 N  R- Q) ^3 E0 ~& o$ E7 `
  35. $worker = new Worker('text://0.0.0.0:2015');
    # \( P4 K1 o# L$ Y
  36. // 4个进程
    9 i6 {6 f, R/ b; q
  37. $worker->count = 4;, q; W0 z  Y) s% G) x
  38. // 每个进程启动后在当前进程新增一个Worker监听
      H9 T* b' e3 J
  39. $worker->onWorkerStart = function($worker)# ^& x8 A' M8 e& c* v5 e. K
  40. {5 t7 Y: T4 T% U7 g; m
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    7 x+ o8 y4 v& N) C! A3 j7 p
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    7 f2 m0 P0 S2 I2 j  s' k; G- {
  43.     $inner_worker->reusePort = true;
    ) y$ W  Q5 u" ]$ O+ @/ H
  44.     $inner_worker->onMessage = 'on_message';4 Z: k* m! c! K' Z& ~# a
  45.     // 执行监听。正常监听不会报错
    / q0 l2 T. X* d3 x# |
  46.     $inner_worker->listen();6 F' k* I4 j2 L& o
  47. };
    7 K7 d: r9 b- k$ s1 Q
  48. ' Q: O8 V% ]/ l- i# n% {) V
  49. $worker->onMessage = 'on_message';
    " D; D4 P9 Z. i8 D, D; M6 {6 a) l' ~1 X
  50. ( X. A# i' Q0 F" {; C, g; K
  51. function on_message($connection, $data)' l7 w) X- u" m8 h' h# O
  52. {
    ! u; q9 S; \! b
  53.     $connection->send("hello\n");
    & Z# i7 b9 [8 h8 E4 C( B
  54. }
    6 e3 I  G) I& [' ~: F) \
  55. 6 [2 f) e# q* b4 r7 s& s' R  A# c
  56. // 运行worker
    - r( U/ b$ i( L  M1 c: E
  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: i, n/ s9 n' \& i) e6 O
  2. use Workerman\Worker;
    7 _) ]6 A% I0 y8 W& C! z
  3. require_once './Workerman/Autoloader.php';
      e' H1 G2 s. {) `! w: r2 Y
  4. // 初始化一个worker容器,监听1234端口
    ; d+ z4 P; M5 V/ `/ Q
  5. $worker = new Worker('websocket://0.0.0.0:1234');+ o$ v: c- D3 I) ^3 y. G

  6. 4 `1 L9 Z! T$ K4 Y: d
  7. /*
    ) ^. [6 I7 v- T
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    ; [9 K0 U  J7 o2 p3 ~
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)- b% E  E" I( P) g, P
  10. */& y6 E% I6 T* S% G) J
  11. $worker->count = 1;
    & G9 \5 j! ^. N" R0 t  g4 ?" H3 N
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口0 L  @( K& k, a% ^* W( Q4 s/ J
  13. $worker->onWorkerStart = function($worker)
    ; g8 u  u, i5 E& \# D
  14. {
    ' P' J. x# I% n! x/ w$ m
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    - c5 R) u! k1 l& Y( [
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');9 _1 i5 c) h, a7 p2 }6 }
  17.     $inner_text_worker->onMessage = function($connection, $buffer)
    1 s% l1 ^+ g8 f! S5 B+ P* \! Z
  18.     {- ]! v- w) d/ @: u6 Y* g8 S
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据
    . s8 ?) h7 \$ m, c) G" i3 |! j
  20.         $data = json_decode($buffer, true);- X& I' R& O& y! v& E2 U
  21.         $uid = $data['uid'];+ I! l% K: Y3 `/ {
  22.         // 通过workerman,向uid的页面推送数据
    % E* W0 F6 W3 Z- d9 v$ y: P
  23.         $ret = sendMessageByUid($uid, $buffer);! j+ U" M* V) }9 q5 K4 `: {
  24.         // 返回推送结果0 r, {  t- x1 |* p' n$ ?" h  N
  25.         $connection->send($ret ? 'ok' : 'fail');
    # N) c- p  b# a4 [9 z6 w% G$ ^
  26.     };
    9 B$ z: H& c) I% I
  27.     // ## 执行监听 ##
    , ]; [: E, f, T0 d
  28.     $inner_text_worker->listen();
    % ]' S$ P, x) n) }0 M
  29. };6 g9 [) D+ G8 G- U# u6 _
  30. // 新增加一个属性,用来保存uid到connection的映射
    2 l3 ?" S' C6 [- S3 m( N. L* g
  31. $worker->uidConnections = array();
    ! J: i5 d: w/ \4 O  w- i
  32. // 当有客户端发来消息时执行的回调函数( U0 k# k. r# L
  33. $worker->onMessage = function($connection, $data)
    & [. n3 [% p& |% g5 y) v5 f* ~: G
  34. {
    7 ]1 Z0 P) v5 l6 `3 D" t
  35.     global $worker;
    . e2 n! }% Q) L5 U
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    , k* O1 p+ V3 b" ?% V, w! Q9 i2 b
  37.     if(!isset($connection->uid))' @( k5 B) _' _$ {6 @7 x% U
  38.     {
    * P- @9 _4 U0 Z3 v+ v+ G
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
    2 M0 W8 h7 M" j
  40.        $connection->uid = $data;* D8 N9 ~* a: Y) ~) I& i
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    9 b  G8 Z) U- L1 }  }
  42.         * 实现针对特定uid推送数据
    " N( C6 n2 d( ]7 l  h& x2 ^
  43.         */
    1 S' Z5 f# T8 y! h$ }6 I: h
  44.        $worker->uidConnections[$connection->uid] = $connection;
    $ {& t  {. r/ m8 {, @' d
  45.        return;
    ! i8 W# P" Y' g5 n; M+ p2 f5 y5 m
  46.     }: J) A- y5 y4 \0 ?5 T) H2 C
  47. };
    ; ]2 J2 b  q  w6 b# }0 g1 t1 \

  48. 8 Z2 E/ \% b+ R
  49. // 当有客户端连接断开时
    4 `+ [3 i+ S, E
  50. $worker->onClose = function($connection)7 ?+ H$ n. V& k1 ?  `1 O: x9 y2 @- m8 P, h
  51. {7 E4 o5 X0 W- d/ d1 ^
  52.     global $worker;
    * R- Q9 W, `" i0 B
  53.     if(isset($connection->uid))) o# F* x5 a$ U0 e6 G
  54.     {/ P2 z) I& M# b" ?3 A9 S6 z
  55.         // 连接断开时删除映射2 z2 e, \/ P9 I4 }. }
  56.         unset($worker->uidConnections[$connection->uid]);
    , U  E' f* G1 C% m( r" q* P
  57.     }
    ( T8 b* R8 H/ `1 ?& V( F3 u/ z! r
  58. };
      G+ W. i5 J0 V" y1 g- d, u* h
  59. 4 D7 j9 D- g* i# O1 r
  60. // 向所有验证的用户推送数据+ G  V* m5 v3 Z& g0 b
  61. function broadcast($message)
    3 u# p1 i2 }+ C" C) @
  62. {2 b0 X# S! k, k
  63.    global $worker;
    , f, U# W1 U% Z
  64.    foreach($worker->uidConnections as $connection)4 {, c1 t9 y& s( r8 d9 ]! c
  65.    {) z7 ]/ X  c! |0 B! o- P1 o
  66.         $connection->send($message);# t- C9 N: S5 @! q. @" U+ T/ }
  67.    }  v7 v) i( j  R! x2 k, N" l
  68. }6 v; ]8 X% d+ `/ i0 ]
  69. ' i' Q6 d4 p5 H" A
  70. // 针对uid推送数据' ~# P- m- G" k
  71. function sendMessageByUid($uid, $message)3 ^7 R! N# h+ h: h; G- F1 @5 F6 u
  72. {
    * `0 N. w$ j# Q
  73.     global $worker;
    ; n% p1 r- G  A" D8 }; F" k; d" }
  74.     if(isset($worker->uidConnections[$uid]))! C8 P8 I: g- z# Z
  75.     {; i7 N) C8 w  d6 B$ u+ Y
  76.         $connection = $worker->uidConnections[$uid];
    $ Q1 n7 N- k: u$ s0 D7 M- H* _
  77.         $connection->send($message);
    ( r, g9 ^! |  {. `- {, t4 _4 @
  78.         return true;
      F; P- {& g/ [3 C6 m4 x
  79.     }
    0 [- @: U# Z. X3 X5 ~$ }4 ?4 H
  80.     return false;/ U2 i! R+ \" a6 ]8 V; F
  81. }& v6 x  I/ p- F/ `; k

  82. + L% m! p# |3 S- v6 n9 X
  83. // 运行所有的worker
    ' l1 r: E* N- e3 p+ k0 ]2 w0 a
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');
    . _2 a  ^2 [" _
  2. ws.onopen = function(){
    ' E3 t: R3 [- O$ }( [
  3.     var uid = 'uid1';
    / q' y. ?& {4 b
  4.     ws.send(uid);
    9 w+ t9 S' s$ W4 n* m6 K
  5. };
    : d( G+ g( z5 P  S3 y9 X
  6. ws.onmessage = function(e){2 M" W5 D; Y5 S/ ^  }
  7.     alert(e.data);
    , A: J/ p8 Z5 L6 ~  J
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口, \$ q/ ^4 A% }- E2 f  B9 z6 v; h
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
    % @" L. G9 n9 q3 G/ D
  3. // 推送的数据,包含uid字段,表示是给这个uid推送- h' O' |, {; ?  A
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');
    / D6 B# }# v4 }, T1 I* g
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符" k$ m1 Q# b) y) g* r
  6. fwrite($client, json_encode($data)."\n");
    . F& ?6 d& s3 g9 N2 A
  7. // 读取推送结果
    0 Z% K. f4 r% v! V, |  G* g* l& {9 |
  8. echo fread($client, 8192);
复制代码
& m5 R: P7 s3 V: N6 i- {7 g

( F7 ~0 j0 N3 e- C  r  w- _! B
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 11:46 , Processed in 0.052472 second(s), 20 queries .

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