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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15614|回复: 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;
    # |  O- [0 m+ \, U+ H3 Q
  2. require_once __DIR__ . '/Workerman/Autoloader.php';9 ^  R* O) ^0 W5 A) ^! y( E4 c4 l

  3. $ K/ r! I) s# u# K
  4. $worker = new Worker();
    . c9 h. C- b# ]$ s0 O: J
  5. // 4个进程; ?9 @- c7 s( x3 N8 B
  6. $worker->count = 4;
    6 v5 H' y* H& W( g4 R( S1 Q# m
  7. // 每个进程启动后在当前进程新增一个Worker监听
    8 m% S2 d; t/ d% J
  8. $worker->onWorkerStart = function($worker)
    , a; b5 ^, q% P: c% D, {9 ]* b- U  ?
  9. {
    6 d: `# L* ~* U% `) W) b9 |
  10.     /**
    . D- c& y2 b6 W& v  Y
  11.      * 4个进程启动的时候都创建2016端口的Worker: }* P4 ?, h( y* {, s# J+ T
  12.      * 当执行到worker->listen()时会报Address already in use错误" a+ J& Y$ z% V9 r3 f
  13.      * 如果worker->count=1则不会报错& z# t0 L" n+ c8 Q7 f
  14.      */
    , R% S% {- k% b# }9 O! R9 z
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');* g0 E. g+ v+ o9 S0 G2 T
  16.     $inner_worker->onMessage = 'on_message';
    % [. `& h" S1 U
  17.     // 执行监听。这里会报Address already in use错误8 J& Y; [3 D) i8 Y5 g- O0 M8 @
  18.     $inner_worker->listen();) E6 U% u" y' F. v
  19. };1 v9 d3 {3 E# Q! Q9 b3 Q7 T3 c

  20. ! {: E" T; _! w2 _
  21. $worker->onMessage = 'on_message';' }3 j% }% f+ r/ A( J
  22. 5 S( A, B* n# h1 r' {
  23. function on_message($connection, $data)
    4 B! K( W$ w6 a4 ], W3 i% C0 q
  24. {8 F$ L  v: x4 _* {% `
  25.     $connection->send("hello\n");% W; J4 p) F7 ~* o
  26. }
    8 B  D% y* K. Y7 t, d- R
  27. : b- [0 ^8 W. f  J$ A
  28. // 运行worker
    $ C+ J0 Y" H* f  A9 T( t8 d  Z
  29. Worker::runAll();
    + A' Z' O3 a" G  A: f
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:* Z2 R5 e  ^; i1 B* @+ a! T5 ^1 E

  31. 7 a% q1 ]% q7 O, I  M) z% B
  32. use Workerman\Worker;
    & H% @  l  `6 N% }$ T8 {& J
  33. require_once './Workerman/Autoloader.php';' Z4 W6 g- L6 [' ^* ?+ O

  34. ' h; k* v( w, D" n1 g, }
  35. $worker = new Worker('text://0.0.0.0:2015');& L5 D) V0 I5 k& u$ X0 q$ k
  36. // 4个进程
    , l% H2 \) f# r6 T
  37. $worker->count = 4;
    7 f( H/ o& E- _4 u) M0 I' J4 D! r
  38. // 每个进程启动后在当前进程新增一个Worker监听5 k: b: [8 [% g
  39. $worker->onWorkerStart = function($worker)4 d$ r( \2 L' F& W0 q
  40. {2 J; D/ \; E' j% s
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    " B9 y$ R6 K$ \" @9 X' J
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    2 k) Q8 C% _! W7 _2 [
  43.     $inner_worker->reusePort = true;
    9 g) O' D6 W! t- R( r) |( R
  44.     $inner_worker->onMessage = 'on_message';
    ; u' X2 A* a1 q# l6 M
  45.     // 执行监听。正常监听不会报错9 G, C! X7 x" D( `. a% N
  46.     $inner_worker->listen();
    . q2 c9 X4 }5 N
  47. };
    6 x+ a- D. D/ B7 L* s" v+ A

  48. + l8 o" N1 A3 d- X
  49. $worker->onMessage = 'on_message';# B, {% P! d9 Z: F
  50. ' `" |. v1 S/ b+ {3 r2 G3 c
  51. function on_message($connection, $data)
    ' H6 e  u* l3 G& E6 O( N6 s
  52. {
    9 S7 |8 ~9 ^- |  n6 Z  H
  53.     $connection->send("hello\n");; b$ C  n- T" d: r; @5 w2 z" A* I
  54. }4 X' ?8 z& F$ m
  55. * |. r2 a0 K2 {! D) w
  56. // 运行worker. g' V" X' a* 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" U! _: ^( U- k9 U/ C$ T/ u/ Z& G$ y. K
  2. use Workerman\Worker;
    ) x. M) y* h3 q: m4 `
  3. require_once './Workerman/Autoloader.php';
    / p" i0 E( v9 |5 Y" I# Q9 c
  4. // 初始化一个worker容器,监听1234端口% v0 t3 v6 [9 u7 u6 n
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    & T1 d: A. S' ?, b- [8 F2 @
  6. ) T5 k! g. Y0 L+ B  U
  7. /*
    + `; x' J6 H" y! l
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误  {4 w& e" q/ A2 g) h; u. n
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
    , \3 j. C* i: @8 ]% c
  10. */
    4 y6 F' _8 d& i  e: x
  11. $worker->count = 1;, O2 Y" C1 c- Y9 e8 y8 ?
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    2 Y/ J+ {: W4 S# o/ ?/ ]5 H" w8 O
  13. $worker->onWorkerStart = function($worker)6 u% Z: H+ w3 R! I' S# I
  14. {* ~# e+ N5 l& U! [+ r
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符, Z& ]/ b3 p0 T9 _
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');$ V& H4 ^$ C  c' q4 g2 p/ |
  17.     $inner_text_worker->onMessage = function($connection, $buffer)  F8 K6 i  T( E8 ?, I( {( k1 R+ @
  18.     {7 @  n6 V( P6 o0 \
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据  H1 V8 I% F0 F: t8 q$ y
  20.         $data = json_decode($buffer, true);4 G2 ~; l; H6 H( Z2 _3 ^% k
  21.         $uid = $data['uid'];
    ' l; I+ j* w4 S
  22.         // 通过workerman,向uid的页面推送数据1 E+ d) f' Z6 P% z1 n
  23.         $ret = sendMessageByUid($uid, $buffer);
    * E* Z. ~2 I1 I9 p3 E1 J& o4 J( J- h
  24.         // 返回推送结果' u# A* |) ?1 Z8 S1 L5 n
  25.         $connection->send($ret ? 'ok' : 'fail');
    - {$ q- b+ ^  ^7 m8 T
  26.     };4 j0 d- A1 x8 F0 r
  27.     // ## 执行监听 ##: N  ?3 w9 D( K* N& |
  28.     $inner_text_worker->listen();" D* E  N: a. k8 K7 R0 Q
  29. };
    ) U" w+ M/ {! J2 w& X' A7 d+ @
  30. // 新增加一个属性,用来保存uid到connection的映射. d& p6 P: F$ F8 i1 U$ \
  31. $worker->uidConnections = array();
    , p! n! W* [8 H0 G' ^& c, f$ A
  32. // 当有客户端发来消息时执行的回调函数4 o0 Q- Z9 H) Z" \4 e) F+ f
  33. $worker->onMessage = function($connection, $data)6 |0 x2 r$ T6 o) q
  34. {% y4 ^1 z0 x1 N- b
  35.     global $worker;3 B) A4 r& j; N, ^5 f% I
  36.     // 判断当前客户端是否已经验证,既是否设置了uid# O. D7 P/ m% m" s% Q
  37.     if(!isset($connection->uid))0 U8 D! G+ e' G( Q6 ~: y5 w& G8 v
  38.     {
    / I/ a( D5 v: w4 ]4 m9 i4 }
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)7 g8 }( Q6 o$ \0 N- I. c2 w
  40.        $connection->uid = $data;
    ' _2 V8 d8 {. O4 _
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    4 S' s) i9 @  k/ t! e% A/ W
  42.         * 实现针对特定uid推送数据
    4 f  ~% ]8 Y7 D
  43.         */! V4 E% ~0 z8 l, n# Q
  44.        $worker->uidConnections[$connection->uid] = $connection;' h! A/ ~- {3 T: V2 V0 e3 B
  45.        return;1 g) W, x" {2 r7 T8 O2 g
  46.     }
    ! m* a: R+ {3 @( I3 w$ V4 L2 x- V7 W
  47. };
    ' G% Y3 \- F( J. @" Z6 H
  48. + ~: ~6 o6 d3 Q4 N9 Q/ f/ c
  49. // 当有客户端连接断开时$ m3 \, h$ v8 }5 a! N- o
  50. $worker->onClose = function($connection)# y" l! i1 |7 B. Q" d
  51. {
    : X: ^4 r+ e1 P  h: v
  52.     global $worker;. {6 g6 X$ K% S) z  \3 ~& k% u
  53.     if(isset($connection->uid))
    ( w) X$ W6 I/ T- `' X$ t
  54.     {4 K1 u3 Z# A  W4 }0 g/ Z8 x$ w3 S* F9 E# m
  55.         // 连接断开时删除映射
    3 q2 Z5 K: u* l( N' [, V
  56.         unset($worker->uidConnections[$connection->uid]);, N) }7 w9 `% |+ {. A
  57.     }
    7 [9 X# v3 {3 l4 M; b: O# y
  58. };
    + E: ]. H$ R, w$ P' E' K( _

  59. 3 S6 c6 t5 }9 |8 T! v7 e9 H' j3 B
  60. // 向所有验证的用户推送数据
    4 H, S8 `* n  Y8 u
  61. function broadcast($message)
    % b5 b' H3 ?7 y4 B) Z' g
  62. {' S5 T( D# r0 Z& l
  63.    global $worker;
    8 ]) ]9 G. r) q
  64.    foreach($worker->uidConnections as $connection)
    & n! O) A6 ^7 O  q1 U8 s
  65.    {
    , E7 M6 \1 f3 y6 Q% [9 ]
  66.         $connection->send($message);6 U+ h! F2 p# ~  d
  67.    }, k# h  A9 r0 }/ h7 R
  68. }
    ' n8 _8 l" U8 m( _2 c7 M3 u

  69. & [, j$ ]4 Z  Z  O4 Z6 T
  70. // 针对uid推送数据7 G, l4 _& R/ P" b, u
  71. function sendMessageByUid($uid, $message)
      L% p2 |& Q4 O" ~( k4 c# a4 G
  72. {, R1 l) n: U, h- M0 c( [# n; G
  73.     global $worker;! t6 G7 z' Y. O- C4 Y- }* N
  74.     if(isset($worker->uidConnections[$uid]))( g/ N3 {, l0 N$ w; k/ ~5 {: ], @2 a' `
  75.     {
    7 ^. D$ U9 B& V5 \% b1 _' s
  76.         $connection = $worker->uidConnections[$uid];
    " \' U! a* C0 A" Y$ G. o  j; D. A5 U
  77.         $connection->send($message);9 Y9 C% }+ T% o( {% W, B- A7 P' m& i
  78.         return true;  v& V6 \7 o, S4 V2 A
  79.     }
    ' q% h8 C/ @8 O. y
  80.     return false;& f  k. _. |1 B. u7 F# n3 v
  81. }0 @& ^- K, {: b) z2 Z

  82. 5 B4 b" @  K: o  O% ]
  83. // 运行所有的worker
    6 k$ e9 Y8 t" ~2 P$ U: R/ T
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');
    ) U" [, U8 E) H+ k( g
  2. ws.onopen = function(){/ m  ?( ^+ _! L9 x( I" y& {
  3.     var uid = 'uid1';2 h, G' q0 g" N
  4.     ws.send(uid);
    6 c7 v  r" f& E9 y, l0 p
  5. };' o+ L. d, b- N3 `
  6. ws.onmessage = function(e){
      ^( o8 n7 a2 b$ |) }) h
  7.     alert(e.data);
    2 u: ^2 S, C% I; c5 W- I
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口, v; Y  U9 p; g# `
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
    / V# r5 P/ I- A$ \+ p' N
  3. // 推送的数据,包含uid字段,表示是给这个uid推送6 J3 w2 _% Z0 c. p1 J5 k  k' H; W5 U/ e
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');  [/ k" @5 Y* L: g8 Z% h/ |
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符1 `* p" C/ q8 {+ Q% w
  6. fwrite($client, json_encode($data)."\n");
    8 j  R0 F& e3 ]0 ?
  7. // 读取推送结果
    + H9 `9 F+ O1 T& Q' v; @) q
  8. echo fread($client, 8192);
复制代码

; Y8 j& K7 z. d) F- g7 r$ }
5 J8 _0 V1 c: A
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 12:39 , Processed in 0.059745 second(s), 21 queries .

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