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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 16160|回复: 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;
    ) x: a: X/ _8 w3 J
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    ! C/ b) S! [# j( `% b- V
  3. 1 ~+ e3 c2 [$ `; }
  4. $worker = new Worker();# F9 z( I4 z. _& R7 G
  5. // 4个进程$ L2 O" d* e. [
  6. $worker->count = 4;& ?- g* P& U5 W3 _2 m
  7. // 每个进程启动后在当前进程新增一个Worker监听
    ) E$ d1 z( k- o
  8. $worker->onWorkerStart = function($worker)
    & h  F1 w/ G2 `' O- @; o0 w
  9. {* ]1 o) O. n9 W# A' @
  10.     /**
    ! y- O8 h- O9 W" w( l6 f
  11.      * 4个进程启动的时候都创建2016端口的Worker
    9 H# a9 N6 q- D* V7 g  k8 w3 B
  12.      * 当执行到worker->listen()时会报Address already in use错误# W# Z8 \8 K( K7 E- h6 s# h8 X7 z
  13.      * 如果worker->count=1则不会报错( v- H0 a* J" n
  14.      *// J8 o2 B* o! T* ?' U! [
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');& C. I+ m4 Y" l, |
  16.     $inner_worker->onMessage = 'on_message';
    ( ~2 x- q  P  C) r5 e
  17.     // 执行监听。这里会报Address already in use错误9 A1 L" j/ G+ v) S; t
  18.     $inner_worker->listen();
    & d# P# }& B0 }4 l% g
  19. };
    " N8 @  q2 E2 I* e3 [9 T
  20. 4 i% |. X; V0 E. F
  21. $worker->onMessage = 'on_message';
    $ R* g# a; y: A2 a2 I. p
  22. 6 H  Q$ L  G' J. x& k
  23. function on_message($connection, $data)$ v' c8 Z/ C4 c  m4 f
  24. {+ r- S8 M8 J% D; l- I) k
  25.     $connection->send("hello\n");% y0 p, e, F( i7 _: N2 d
  26. }
    ' c- j/ l3 A: H: A

  27. ' V7 D: K/ r! w. x
  28. // 运行worker
    1 I) k2 \& J* V: F) e3 G
  29. Worker::runAll();8 N, z2 _5 @( s3 \% F8 C- S# U  a
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
    ( O5 B, {) H: p
  31. . t4 [! S" X# |+ ]1 g8 o, R$ E
  32. use Workerman\Worker;& ]( h: i  ]4 A0 V. B6 T
  33. require_once './Workerman/Autoloader.php';( o, q1 Q; x. z' U8 b' Z# y

  34. - Z) q  P. }( K! L8 D5 N
  35. $worker = new Worker('text://0.0.0.0:2015');
    , F& }/ R- R9 b  w& O* U" e
  36. // 4个进程
    7 R$ W% ?6 ^! c1 R) K1 Z- d
  37. $worker->count = 4;
    ( B3 S7 U2 n, f) N& ]
  38. // 每个进程启动后在当前进程新增一个Worker监听: Q' I6 d$ j8 \6 D3 n9 R8 k; F) E
  39. $worker->onWorkerStart = function($worker)
    2 V% F: a% r4 s1 F5 R
  40. {$ g/ g* e; o! d( ^
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    1 I& \8 G) y( W! C; Z% {; B* h
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)$ ~! {0 L5 s5 M3 l
  43.     $inner_worker->reusePort = true;: X3 R! z1 U* ^# s6 Y
  44.     $inner_worker->onMessage = 'on_message';' F& {7 O* f* u/ X
  45.     // 执行监听。正常监听不会报错
    . E1 R$ ?, x+ u
  46.     $inner_worker->listen();
    5 m; u. Q) F. J/ g: ~
  47. };
    : |, {0 K$ J$ o6 E# M

  48. % q: K+ A6 o0 D1 t$ z
  49. $worker->onMessage = 'on_message';  y1 T5 ~, t$ T3 x- t

  50. 0 C3 T9 Y5 ^# c% N2 z8 q" ?  N& C% L
  51. function on_message($connection, $data)+ ~% G- d( S: T9 d: u/ z8 Y' Q9 J
  52. {" Q8 P: q" `9 A' u" u& R
  53.     $connection->send("hello\n");
    8 j) w8 ~0 D- R) j) G: X. {+ P
  54. }
    5 _  x/ f& c$ H) R
  55. + B6 C, {8 t# a, X, t
  56. // 运行worker* N; H2 X6 w2 @. A
  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. <?php9 S0 w) P+ j6 W+ k$ n/ p0 p2 [) T
  2. use Workerman\Worker;
    / J8 v! i: V! [. u
  3. require_once './Workerman/Autoloader.php';
    4 n% I, o% u5 r* r6 h& E: n8 m
  4. // 初始化一个worker容器,监听1234端口
    / g4 z( \! @$ I
  5. $worker = new Worker('websocket://0.0.0.0:1234');
    * f5 P9 e2 j  k

  6. 8 U; s' n+ e9 ]
  7. /** A1 i/ g6 k2 ]5 D& ^, s  h- c" Y8 V- T
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误! F& L/ i' x; j# _4 D+ ?$ j
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
    , c  y! X0 l/ S2 i% o% z3 M& \% n
  10. */
    ! L6 A' c2 x; v( y9 V9 T
  11. $worker->count = 1;
    $ t, z/ ^$ o7 e( T3 I8 Y
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    6 b! c5 K5 e8 @1 P6 T7 u1 m
  13. $worker->onWorkerStart = function($worker)
    . a5 I% `7 p: f0 U7 D+ w3 Q
  14. {& ]7 c8 l/ [3 D- \
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符$ a& ]$ c1 |* b/ Q
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');5 t1 c+ t/ I9 b  E) T% H; N5 J" V
  17.     $inner_text_worker->onMessage = function($connection, $buffer)/ R" ?/ E5 }4 o! T$ W  _/ s
  18.     {
    ; @/ v$ E7 W8 R# b' G
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据! U+ h7 f) m! G% V+ n1 w
  20.         $data = json_decode($buffer, true);
    5 I8 t5 ~8 M  }" }: q2 V9 X6 A/ n
  21.         $uid = $data['uid'];
    8 P: X1 S/ g4 V. Z6 g/ s2 U" G
  22.         // 通过workerman,向uid的页面推送数据
    - M: t( G/ H  F2 g) p
  23.         $ret = sendMessageByUid($uid, $buffer);3 `4 h1 D: ?# B/ C$ r( s
  24.         // 返回推送结果
    & A, Q7 l7 ]( N' S
  25.         $connection->send($ret ? 'ok' : 'fail');0 l3 U1 c7 N, Z
  26.     };
    + k4 X) F" c2 S4 K- e
  27.     // ## 执行监听 ##
      `& X5 T- T6 K* I1 |) d
  28.     $inner_text_worker->listen();
    5 l$ P5 u. _- o+ v3 D: F! c- p
  29. };5 v. j8 K9 k1 W8 V; l
  30. // 新增加一个属性,用来保存uid到connection的映射" }  |' s- K! z# D+ G
  31. $worker->uidConnections = array();
    9 K; P$ s2 y- }: F) q
  32. // 当有客户端发来消息时执行的回调函数
    , [( B5 f8 K$ @2 H  I# W
  33. $worker->onMessage = function($connection, $data)) @; y- J; o) n# X4 Z) {; g
  34. {' _  T' y% T) t1 v, q0 H9 L% g' D
  35.     global $worker;  b- h9 B* F( K- I' ^
  36.     // 判断当前客户端是否已经验证,既是否设置了uid
    ( S9 B" o! D& r% B/ X& H' J- ~/ ?
  37.     if(!isset($connection->uid))
    5 I; N* U% f9 G- ^  f$ H0 z( P) D
  38.     {! R7 J. |; ^( }
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
    ) A# k/ G, ^7 w/ n
  40.        $connection->uid = $data;
    & Z: g0 K& t# w) A& t
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    - b* z3 [$ G" O( \8 d5 _
  42.         * 实现针对特定uid推送数据
    3 X4 K" h+ J; O" \1 m. D6 B  q% j
  43.         */5 B' }+ d) W9 R* y% M
  44.        $worker->uidConnections[$connection->uid] = $connection;) ~( P# U; `: Q1 O  T$ a1 x8 f
  45.        return;
    , C  R4 @9 S7 @$ W9 ?
  46.     }& T/ B2 t9 v6 E# `$ S7 O0 ?. {
  47. };! f6 y6 U2 |5 X4 d

  48. 8 n  V. w& [% J  j- j' r7 P3 H+ F" B
  49. // 当有客户端连接断开时
    8 F1 @1 m0 X& r6 L
  50. $worker->onClose = function($connection)" }7 C- L  |) k" X* H
  51. {  B" P% L# k# j
  52.     global $worker;+ h1 C; @! h+ q) d6 p0 O1 X
  53.     if(isset($connection->uid))0 O; \6 J6 F. r, U7 k. s" M5 E
  54.     {
    / ]; t% @5 X% y0 G' f/ O5 I* G7 V
  55.         // 连接断开时删除映射
    ( B! Y1 B$ }/ I3 Z6 I; q2 x
  56.         unset($worker->uidConnections[$connection->uid]);
    5 |1 P: G# B  j! G0 L; q9 G
  57.     }
    2 L8 y( c2 c" D
  58. };
    , z' F1 }( t! I" t5 H

  59. " n' p: {9 V+ d1 h$ _( ?0 `, r
  60. // 向所有验证的用户推送数据
    4 a' I  R& v0 s. B4 I$ }
  61. function broadcast($message)% i7 D. g# A* V  T. h* B7 m
  62. {. Q  m& M( C9 ?- b% i
  63.    global $worker;
    " c) Y: L2 I, }* H) A, r  j
  64.    foreach($worker->uidConnections as $connection)
    % \( ]2 T6 D  v& A$ @8 O# K
  65.    {
    9 K; r. y$ x& |' n! |
  66.         $connection->send($message);
    4 Z  f" U. j, K; g
  67.    }  z1 o8 j! m2 U, F6 q
  68. }- t3 Q4 k/ U6 s/ V1 q/ p

  69. ' H5 @% j, K+ t" s3 }
  70. // 针对uid推送数据6 E2 T3 ?; h6 L- F# C/ K5 q' k$ r
  71. function sendMessageByUid($uid, $message)( A' P6 r! q  N: @
  72. {
    . I3 l& \2 H- m0 s! l- f
  73.     global $worker;
    * c2 w* Q; d0 q
  74.     if(isset($worker->uidConnections[$uid]))
    * S3 Z5 y7 j% V
  75.     {
      N5 N) `% |% u9 U+ D+ l3 Y
  76.         $connection = $worker->uidConnections[$uid];
    1 {  R' ~/ {6 d& X- P0 U
  77.         $connection->send($message);- T3 i6 |0 D5 C1 ^4 `
  78.         return true;" h5 h9 {& P! c, ^
  79.     }( H2 P" _' Z7 j
  80.     return false;
    1 l0 m' b* \( r3 |0 z
  81. }
    6 s+ h3 b1 d3 I# ~2 P! B; O5 A
  82. * ?/ I$ n3 g' u9 O' n
  83. // 运行所有的worker5 d: k2 @  i9 @
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');& l9 k4 ~  n! E! e# ~
  2. ws.onopen = function(){
    : O  h3 W1 v6 r) ~
  3.     var uid = 'uid1';5 R8 s8 [# x9 x  n
  4.     ws.send(uid);
    . |7 y2 }$ p  |* m7 j  \7 Z2 v
  5. };7 {" q8 N' z( c) Y" X
  6. ws.onmessage = function(e){
    * n% Y+ X3 D+ O1 P, k, E* s
  7.     alert(e.data);3 Y4 q: w+ W4 A% e6 F
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口  U: g6 y- s* x! t+ o. {% a
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
    * C+ ~2 A& I+ \- X$ W3 ]
  3. // 推送的数据,包含uid字段,表示是给这个uid推送6 ]/ j+ _) w  ]4 E9 X- K; }9 O9 E, a
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');5 N3 N- J2 Y( D8 `8 }9 Z* R4 k6 s' |3 b7 D
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    . o$ p% I4 B6 Y) m6 h
  6. fwrite($client, json_encode($data)."\n");. [, w$ p2 T9 D. J; y6 W" \
  7. // 读取推送结果4 |* e" v, ^" j, w1 }2 a9 `' Y
  8. echo fread($client, 8192);
复制代码

: W, q7 C% e- u
! G* w" ]' o- q
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-9-21 03:12 , Processed in 0.059938 second(s), 19 queries .

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