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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 16188|回复: 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;# Q8 c1 a. X3 V, C  M7 N" |1 M, P
  2. require_once __DIR__ . '/Workerman/Autoloader.php';0 X* ~$ O: B& R& [% f

  3. " @; X7 ]0 M8 v
  4. $worker = new Worker();
    / E& {: I5 y2 \% O, E$ J6 @
  5. // 4个进程8 h; R* q3 f4 M7 p/ z
  6. $worker->count = 4;+ u0 j3 T$ ?2 c$ ?0 i$ \  ?
  7. // 每个进程启动后在当前进程新增一个Worker监听) s3 a, f4 C* O, v) t' q7 `/ g
  8. $worker->onWorkerStart = function($worker)- G! ?7 j' D1 ^; f, d2 `! t
  9. {, P9 {, @/ \1 R. h# w8 S
  10.     /**
    0 ]' K2 K1 T' a! c# b  n
  11.      * 4个进程启动的时候都创建2016端口的Worker
    8 [/ n# [) ?# h3 G* |$ ^
  12.      * 当执行到worker->listen()时会报Address already in use错误
    0 c; t/ [+ I4 l: S( R1 v6 m
  13.      * 如果worker->count=1则不会报错; r) \. t- h* d2 _4 f2 }4 {9 J
  14.      */
    6 W! W0 k) N9 P& u2 p7 \
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');9 K- x- s. z' \* w6 v
  16.     $inner_worker->onMessage = 'on_message';# g* Z& Q- K* @( B. P
  17.     // 执行监听。这里会报Address already in use错误5 c( I( {2 z6 ]" ~7 S! u
  18.     $inner_worker->listen();
    % u8 j4 P- i  r, p9 u
  19. };( J) K$ U  ]" j6 w8 H; W8 }
  20. 8 i- V- c5 M- B/ m9 z4 W' [
  21. $worker->onMessage = 'on_message';
    8 w' t. C  Z2 X) E

  22. * O/ B. I$ {% |# ~
  23. function on_message($connection, $data)% u5 \4 @3 ~% U2 N8 m5 m3 g9 h) V
  24. {+ ?7 x9 Q& o/ q" {
  25.     $connection->send("hello\n");6 \& p* }6 `  r
  26. }
    8 e& o# a) B: I' ?3 R3 h9 j

  27. % b' N; b. f4 Q5 C& w3 }- d5 E/ T3 l
  28. // 运行worker8 a: S7 V( n* h9 w3 s
  29. Worker::runAll();7 v) {$ w9 r, N2 x7 [5 u2 N" g
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
    9 o* A1 Q+ i" R2 |
  31. # \/ J9 D2 J9 \, |' {
  32. use Workerman\Worker;5 T5 X) [' F- q/ [- ^" |* p. J
  33. require_once './Workerman/Autoloader.php';. D9 G8 R! [) c! a( a" i: R

  34. 6 D1 l; K5 W: U; y1 M
  35. $worker = new Worker('text://0.0.0.0:2015');
    7 L- R& s+ v! n
  36. // 4个进程; q3 k$ T1 ?% G- q
  37. $worker->count = 4;& q7 g" F4 `: k
  38. // 每个进程启动后在当前进程新增一个Worker监听5 D. u7 w1 e! r& P
  39. $worker->onWorkerStart = function($worker)& w' M1 B5 n9 }% q% x5 L; j) {9 ?
  40. {0 B& P% S1 \! t; c
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    3 l9 I$ ~7 Q" L, m, B
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0); z) U3 c; Y% D/ H
  43.     $inner_worker->reusePort = true;. B, T3 W  H  j+ y8 T# r; ?0 C0 ]
  44.     $inner_worker->onMessage = 'on_message';$ \" ^4 x- H. w4 l% x
  45.     // 执行监听。正常监听不会报错& }* w2 M) T# `5 j( p6 z/ g0 A& F
  46.     $inner_worker->listen();
    4 L0 z; N$ V' H# c+ M6 j
  47. };
    9 _# R8 }) ^$ n- `

  48. 8 `3 Q& S+ j2 s: l6 n( R. W
  49. $worker->onMessage = 'on_message';
    ) s/ h  @" A: k* h2 a% @
  50. 8 K7 C; O  x/ @& t1 n2 ?
  51. function on_message($connection, $data)# S$ J1 c. O8 n) t& O/ K, \% b
  52. {( j  @, J+ B' K! O
  53.     $connection->send("hello\n");
    " i% C) @: e0 s4 f! \! t, }
  54. }
    ) `; Q! H& x5 l

  55. % C" T2 W0 F, N
  56. // 运行worker$ D6 O+ l( I( K/ B; M
  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. <?php1 I% X& V" W" K
  2. use Workerman\Worker;7 w# u  S! ^; f- t& V+ g6 j
  3. require_once './Workerman/Autoloader.php';$ U9 `% |4 Z; W( Q$ f: x0 G% D
  4. // 初始化一个worker容器,监听1234端口
    0 L0 m+ K+ d/ Y& V$ O
  5. $worker = new Worker('websocket://0.0.0.0:1234');& v+ y9 u: B; c6 l3 G, h% O4 [8 ~
  6. : z# _# l) a+ u2 P" v/ j
  7. /*, Z6 C+ F* f/ V2 J( D2 Z
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误
    " ?2 }0 ]5 J! Z- w) d  B: ^( V+ k
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)# p  L; `' j$ K1 c
  10. */
    / I8 F/ x" U8 n' S
  11. $worker->count = 1;% A9 E5 W) }+ u+ f3 B6 c! D0 v* ^
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    9 k2 h8 p2 Q; I9 z
  13. $worker->onWorkerStart = function($worker)
    : u4 N4 r! i8 E+ i
  14. {+ G/ c% y, k# L) v+ p& G
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    1 A; H% ~6 h  [7 ?& A# B% e
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');; b/ d( f  B- s( @' q0 l' k
  17.     $inner_text_worker->onMessage = function($connection, $buffer), H- ^: u( v9 a. c1 M
  18.     {2 G' W8 @6 ]% }# {. e. Q' X7 T1 M
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据/ Q& l, g' e/ `5 h
  20.         $data = json_decode($buffer, true);4 @, w2 g6 V% Z+ R/ w& v/ w
  21.         $uid = $data['uid'];6 T9 s: a8 n" W) t% h+ H
  22.         // 通过workerman,向uid的页面推送数据
    " N$ o* i- c4 j' u" k
  23.         $ret = sendMessageByUid($uid, $buffer);/ e: p" c1 [" {. x( P
  24.         // 返回推送结果
    - D/ o, a3 W3 ^& A  @3 m
  25.         $connection->send($ret ? 'ok' : 'fail');
    4 E. H( d, U! L: g; Y( W
  26.     };9 c) y* e; @8 g  \4 r
  27.     // ## 执行监听 ##
    3 Q8 Q' y: X; p4 P4 [0 X1 ]+ G
  28.     $inner_text_worker->listen();( ~# P8 B# a3 r5 ]: [2 e5 }* @
  29. };
    6 A0 m/ ]" `, h
  30. // 新增加一个属性,用来保存uid到connection的映射
    ( i! i% w6 j# Y7 e4 `% O- N. q+ o
  31. $worker->uidConnections = array();
    # x7 ]+ R" h+ W7 q& J/ p
  32. // 当有客户端发来消息时执行的回调函数7 |9 W8 j2 p4 r  f2 q- h# w
  33. $worker->onMessage = function($connection, $data)
    $ E' a" R7 i4 [
  34. {: a2 I/ y* \2 a% B) P1 j3 e
  35.     global $worker;
    6 d% p! u1 c1 a, i4 ?
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    % X% S8 M8 Z: O4 b% ]; i+ g
  37.     if(!isset($connection->uid))- j2 ]: n# e% W; I. k/ g3 m3 R
  38.     {9 T/ P  K8 ]3 j) K
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证): t, o. ^' q- y( N+ f! x
  40.        $connection->uid = $data;
    6 r7 Q$ n% F$ |# R+ z
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    + a2 T# ]4 P# l3 y4 l8 @6 b
  42.         * 实现针对特定uid推送数据
    2 H# Y4 s9 D3 i: b" e0 P
  43.         */( A1 r- I0 W) o" Y. U0 a
  44.        $worker->uidConnections[$connection->uid] = $connection;
    ) _7 m) P6 p- e/ ^, I$ I$ |9 f$ {
  45.        return;
    $ e( v: H: U, s
  46.     }
    8 `, m% h. y5 N! x
  47. };7 K% v" w8 f% ?( b% W6 ]0 G

  48. ' w9 n+ c% a8 E; ~
  49. // 当有客户端连接断开时
    / {* ^+ z/ Q& L: s
  50. $worker->onClose = function($connection)
    + D% p) q9 |. S
  51. {
    + `% l9 g# n6 B# B. O
  52.     global $worker;9 \& d- z! ~% L( J) \
  53.     if(isset($connection->uid))) J. f( B" K- ~$ s( E  i! ?6 m
  54.     {
    5 C' u1 S5 B' u' |# `
  55.         // 连接断开时删除映射- P) _  ^2 _9 f6 E8 k( z( O- l" \
  56.         unset($worker->uidConnections[$connection->uid]);( x6 _; R7 r- d: C
  57.     }) C6 m$ r! n+ o1 v+ e4 B9 y
  58. };
    / B% [* q8 V: Z0 Q( k6 [7 Q* y: b
  59. ' v, k4 S$ {: ]. s  n
  60. // 向所有验证的用户推送数据
    * ~( f0 M- q* j% m  R3 t7 r
  61. function broadcast($message)
    ! a* X6 W; r! A. _5 x" O
  62. {
    2 ^  Y- k' m- r2 l) s
  63.    global $worker;8 @  h" l5 }& A& z
  64.    foreach($worker->uidConnections as $connection)' H3 t  m% c+ L3 L' I7 c" K0 ]
  65.    {$ u. q) D: ?: w  ]
  66.         $connection->send($message);# h: u) n0 u* i# N3 o5 N
  67.    }
    ' Q, q1 v+ U/ P6 y- ~
  68. }
    3 U% u* a+ A2 \; F1 @2 h( J( ]

  69. 8 W2 Z" Q, b0 H, F
  70. // 针对uid推送数据" {0 ^$ r/ K. ?; d; H
  71. function sendMessageByUid($uid, $message)" g, P. i9 n( c
  72. {
    2 U7 k, G! t' e. W* Z' |
  73.     global $worker;
    : X1 `6 W7 n! i/ {: L0 }4 r; l
  74.     if(isset($worker->uidConnections[$uid]))
    % z. H, J' x* W4 J4 s
  75.     {
    & D) E0 O6 c/ N: M' P: u. [0 U
  76.         $connection = $worker->uidConnections[$uid];% q* H/ y) X. \- }
  77.         $connection->send($message);8 {! Z( C3 f# b3 e
  78.         return true;
    . r6 h! Q  b- c$ P& G7 P
  79.     }0 R7 X6 p' W! N, b
  80.     return false;( `: S( i8 U6 K; |( x
  81. }7 }' T- u( P! A

  82. + [/ @4 J! A, B, `" p* X8 N
  83. // 运行所有的worker
    1 M2 H. S4 O) W; I
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');
    ! m9 b2 u! `& b4 l7 q, y; @7 K
  2. ws.onopen = function(){
    ! a# @( Y3 x# x4 r! Z/ `* K9 |
  3.     var uid = 'uid1';
    9 G5 x  y& ^" y! Y5 @$ a5 L  E# Z
  4.     ws.send(uid);
    # m1 U" y9 _0 G- J1 i
  5. };9 y- y. N' n+ `0 y9 ~
  6. ws.onmessage = function(e){
    4 S1 w) S" D% n: y7 ~$ c: L* J6 u+ Z
  7.     alert(e.data);% v2 g9 b0 a" Y+ H  {
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    ; I6 H* [6 v, {; a# M5 M5 H
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);1 `6 a5 i, n/ y" Q  I1 }% }) \, C
  3. // 推送的数据,包含uid字段,表示是给这个uid推送, t, t: U8 h9 T7 y7 d
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');3 p! k; U7 H2 w* V, a
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符+ l* e6 ?6 ]  G
  6. fwrite($client, json_encode($data)."\n");
    ! N  {7 g: z& W& }5 a
  7. // 读取推送结果
      Y8 W( C, K0 l2 ^( d5 k" S
  8. echo fread($client, 8192);
复制代码

0 \8 j+ C& L6 q! G( n( _4 ~4 [
* m+ ^% `$ v& ?7 T, Z7 Q0 g
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-9-22 04:11 , Processed in 0.054063 second(s), 21 queries .

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