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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15617|回复: 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;
    : u6 d+ w8 t# o  b+ x- E& e
  2. require_once __DIR__ . '/Workerman/Autoloader.php';- b$ @/ o: ^4 X
  3. 1 s6 i( e$ t! l3 H" U
  4. $worker = new Worker();
    1 B( ^9 H* \$ U6 X1 O
  5. // 4个进程
    & C8 J; j( ?" s7 ~, s
  6. $worker->count = 4;
    7 Y6 |0 |  u0 s% t* g' L
  7. // 每个进程启动后在当前进程新增一个Worker监听
    ' |3 F4 B$ D7 Z6 Z
  8. $worker->onWorkerStart = function($worker)
    1 B9 O+ h/ L8 c0 I
  9. {
    5 c: `2 U9 S9 c5 \. d, p: w! I
  10.     /**
    7 |# g2 b$ w8 B; ?
  11.      * 4个进程启动的时候都创建2016端口的Worker
    0 u# _$ ]5 r) n; h  r8 Z5 s
  12.      * 当执行到worker->listen()时会报Address already in use错误0 W  A# K  I" C3 t3 o
  13.      * 如果worker->count=1则不会报错
    9 \( s* `+ ?9 b2 E5 p
  14.      */$ r* x/ q% w" O2 X  w
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');
    6 y7 l& p  \& U$ ?/ M
  16.     $inner_worker->onMessage = 'on_message';8 p) F, F9 ]- v* x( x# W
  17.     // 执行监听。这里会报Address already in use错误
    ( ~; u+ P+ X/ ~: ^- c
  18.     $inner_worker->listen();
    7 I; Z' t2 @, t7 ^
  19. };- t3 a; |) A! Y: M

  20. - @* f. U8 Z  V5 z, |/ j
  21. $worker->onMessage = 'on_message';
    : h, H- K+ ^3 O/ z: z3 g

  22. ( Z5 Q& k# E$ u5 L
  23. function on_message($connection, $data)
    1 V# k- q* Z( X' B
  24. {
    " y4 ~# @2 o. k; K$ S" G
  25.     $connection->send("hello\n");# {3 o0 M3 @6 o9 x. V
  26. }
    1 I9 x, R/ I! G4 \1 q5 g4 W

  27. 3 r- _- e. P+ J
  28. // 运行worker
    ) `! g6 s9 Y5 q
  29. Worker::runAll();
    7 u. p3 R! e: |0 a6 U, F
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
    9 n1 a+ y- m# y& x# J
  31. - P. I' ^7 L9 O8 {
  32. use Workerman\Worker;: B8 I% U8 _& o8 G
  33. require_once './Workerman/Autoloader.php';  L  V7 L  t# k; y

  34. 5 \& V( P" e- m# V' ]  @
  35. $worker = new Worker('text://0.0.0.0:2015');6 s7 ^; N( F6 ^' r3 `: ]
  36. // 4个进程
    , s$ J( a; k3 b
  37. $worker->count = 4;
    - ]$ W1 d- R$ V1 ?' l
  38. // 每个进程启动后在当前进程新增一个Worker监听
    7 H: P" `& Z2 Y$ H$ C* s, E
  39. $worker->onWorkerStart = function($worker)( L+ Y% ~0 P; T. J4 u8 p
  40. {
    . l9 R' h/ d1 L& `6 A; l
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    3 P! ]2 P* C) M9 @3 O/ Q
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
      B) y1 w& M' c+ s- A: G
  43.     $inner_worker->reusePort = true;4 o; w, w' Z* \4 {9 t7 V6 q1 V# v4 S
  44.     $inner_worker->onMessage = 'on_message';3 h- M  ?8 j6 ?. V' {  o
  45.     // 执行监听。正常监听不会报错
    ! U2 |& g0 Q0 h: o
  46.     $inner_worker->listen();
    " l- e' F" ~# L3 ?
  47. };
    : b2 f1 _( ~( {$ l5 A, g6 N4 p+ k- H" B
  48. * p5 ?: o- h* K0 F# T8 H
  49. $worker->onMessage = 'on_message';
    % g) y: ?4 r5 E5 B

  50. 0 y3 F% h# ^& n( u9 D! D8 o8 C- e
  51. function on_message($connection, $data)" ~% ~" T5 f+ m# X7 D" m1 G. J5 N# A
  52. {
      B& \+ {9 g' Z3 q/ C5 N
  53.     $connection->send("hello\n");- P8 ]2 B* V1 }  C, p
  54. }+ t; \& {$ b9 j& ~  z1 Z' s( e+ y+ u
  55. & j& ^! F1 A" r4 f) n, x8 L
  56. // 运行worker' x  s8 A% v( ^% x% l- |+ o
  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. <?php8 h" A+ E4 Z8 Q, U7 g
  2. use Workerman\Worker;9 \4 {$ w) y2 G: Z& [
  3. require_once './Workerman/Autoloader.php';
    , N/ n7 ?2 O/ k7 D& K$ k
  4. // 初始化一个worker容器,监听1234端口/ F" n4 N% s2 `+ o
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    1 w: O2 W# h. ~+ ?* L
  6. " E2 F% B- w5 D8 r
  7. /*
    0 Y: \9 f4 s1 J$ r# n
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误7 s8 r# l& @  F. p; C
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)9 L2 U7 R+ n3 M7 @& u( H( W
  10. */
    + U3 ]1 p+ t2 z6 |
  11. $worker->count = 1;
    , A2 e: u% j% z: |- x2 Q# j( k
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口! s* S( K/ ]6 Y5 q' o
  13. $worker->onWorkerStart = function($worker)1 o* i$ T0 F% P, P: N
  14. {- C& E9 @: J0 b, E/ d
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符3 _- v& A* k* p" ]  w3 T* i2 k- @
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');
    3 A) y0 |2 _' c0 l$ K
  17.     $inner_text_worker->onMessage = function($connection, $buffer)
    $ t/ q  y: V- ~" X& y
  18.     {
    * N/ ]) M! f" j4 ^9 S& P8 M
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据+ n! L2 b) B. T" K6 t8 h9 M
  20.         $data = json_decode($buffer, true);* e0 i. i' a& V4 s; _
  21.         $uid = $data['uid'];
    * c: G0 B5 z+ ^  p. Y: E" t; D( Y
  22.         // 通过workerman,向uid的页面推送数据
    $ k/ ]- `* l/ |# P+ v5 y
  23.         $ret = sendMessageByUid($uid, $buffer);
    ! t; ^) s0 N3 Z8 ~/ A( E: a
  24.         // 返回推送结果2 _7 \. ^2 N5 [7 d; c6 q  ~9 A; a
  25.         $connection->send($ret ? 'ok' : 'fail');
    ' T+ S0 a3 V* M5 Z& [# \  [( z
  26.     };
    6 N9 j7 q: c! p; G/ Q# j0 B. z0 K
  27.     // ## 执行监听 ##. d" n. \8 O( i1 @% G$ ~1 ?
  28.     $inner_text_worker->listen();7 B: ^6 r% N5 k* l4 H8 k2 @7 J$ \
  29. };
    & [5 c  k8 Q9 o1 G+ G9 y
  30. // 新增加一个属性,用来保存uid到connection的映射
    - c  L' q% {' b) R, W8 s8 b
  31. $worker->uidConnections = array();( q( V, c; R! Q: X+ P! L3 [
  32. // 当有客户端发来消息时执行的回调函数
    0 v; A  W+ _- B0 A, C, p
  33. $worker->onMessage = function($connection, $data)& v/ H) o( k- }
  34. {4 r2 }+ J# G; S) }& c
  35.     global $worker;
    1 l$ c! Z1 `4 ]# ^* y8 {  I% P: F
  36.     // 判断当前客户端是否已经验证,既是否设置了uid: D3 a" d" `- p! C
  37.     if(!isset($connection->uid))& X! g6 f& f5 M
  38.     {, p" X: X: q. V' y: r! o
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
    0 P+ G3 g, X% F! J5 H, z
  40.        $connection->uid = $data;
    * J2 p; u9 `, s, m9 |7 f
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
      F/ b8 b5 w1 G/ O
  42.         * 实现针对特定uid推送数据
    - R6 u' Q0 a+ {! S$ q6 o
  43.         */
    * w* Z- k- ]3 C
  44.        $worker->uidConnections[$connection->uid] = $connection;
    6 B5 Z+ w5 h: t! ~# d! r. w' G& V
  45.        return;
    3 ?! q! N) O" t. J
  46.     }2 ~5 i" Z- b: F5 M9 h
  47. };
    # O& Y9 U+ U) C# K  x; |. F, G% j
  48. + {+ z6 }: V. m! K  ]
  49. // 当有客户端连接断开时
    3 _( _) r+ J% }& F1 _5 X
  50. $worker->onClose = function($connection)' N  ^9 @4 R/ J/ i
  51. {
    % |' g; R& |* R' P/ z
  52.     global $worker;
    5 e) `2 _2 O- U" T) }% M
  53.     if(isset($connection->uid))3 S# v9 f0 X; L1 N+ x' u
  54.     {5 _" F, I- W. J5 t+ v
  55.         // 连接断开时删除映射) h% X' C1 T2 [2 d5 z, |
  56.         unset($worker->uidConnections[$connection->uid]);
    ! t. D9 U. n5 o  X
  57.     }2 b( \5 K1 K) ]+ l
  58. };, k( k: N- H1 L- i! g: [
  59. ; p# m6 t4 T0 V' S4 ]3 X
  60. // 向所有验证的用户推送数据
    $ U7 Y1 r) {4 G9 T/ |
  61. function broadcast($message). o* g$ K& R9 C) N8 r* |
  62. {
    4 {# R/ e6 Q1 I. e) T: b# {
  63.    global $worker;, ^1 `" d( h% Y, f( E, L0 l
  64.    foreach($worker->uidConnections as $connection). F( v0 D% h. [" Q
  65.    {
    , h6 W+ c: _+ x+ |) r
  66.         $connection->send($message);
    8 v7 c( d4 K% ^! O4 U
  67.    }# ~  `. |; C6 B. U) h% F
  68. }& E; v: E+ _3 x$ I. J

  69. ! b( X3 e" C& ^2 }# b
  70. // 针对uid推送数据
    ' V& A; `0 J0 j. u9 Z, {
  71. function sendMessageByUid($uid, $message)
    2 V: a! i, r* }& S' L' F
  72. {  B5 x2 J+ B" I" f5 f- R1 k; @, E
  73.     global $worker;( S: s4 r: F7 T! S' c
  74.     if(isset($worker->uidConnections[$uid]))
    * n0 p- d7 o: ]9 A
  75.     {
    8 i/ Z$ F% H3 {* v  j0 Z/ Y
  76.         $connection = $worker->uidConnections[$uid];
    4 R) r6 g, ]6 j0 p$ Z8 d1 A
  77.         $connection->send($message);2 O$ k& L& s+ I8 X& U
  78.         return true;
    - n: r7 D3 U/ A: H6 ?% J
  79.     }( r& R+ Q8 M: D! C( \$ R
  80.     return false;
    ' s8 ~0 _- d! _2 y$ P" M1 p, [
  81. }
    0 {$ U; O, K+ q' m

  82. 6 ?! z; I8 p2 a& x- V) A0 Y0 |
  83. // 运行所有的worker: W2 d* m! G! ~4 a7 G
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');# @& N( O  `7 k; P8 l
  2. ws.onopen = function(){2 Q7 |* _4 A0 p
  3.     var uid = 'uid1';4 v: O- f8 b+ s, |6 Q
  4.     ws.send(uid);
    : \4 V) ]" O8 h( Q% U
  5. };
    , X% L* C/ `8 s  J( f: y7 M6 w
  6. ws.onmessage = function(e){; `1 u; Q, m2 `$ E, t) p9 l7 Q
  7.     alert(e.data);
    * e: Z3 \; K' `8 |$ S# C% e% @3 `( n
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    * p) q' d8 c* Z7 V. E
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);4 V7 p5 R3 p% X" Z
  3. // 推送的数据,包含uid字段,表示是给这个uid推送
    , B. ~: T. w+ ?; u
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');- i- o$ d' N; |7 S  X
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符# {0 i; [7 x$ l" ]4 i/ z
  6. fwrite($client, json_encode($data)."\n");2 I: A4 ]0 t% [0 ], `7 r% y
  7. // 读取推送结果
    % v; R  v% q1 K9 T) j  e4 ^  W
  8. echo fread($client, 8192);
复制代码

4 Q+ Q1 \" [  d+ X# y2 z7 F! G/ ?9 C/ Q; l& f
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 13:13 , Processed in 0.054183 second(s), 20 queries .

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