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

QQ登录

只需一步,快速开始

 找回密码
 立即注册

QQ登录

只需一步,快速开始

查看: 15615|回复: 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;: U7 z& u* ^/ R% o. V
  2. require_once __DIR__ . '/Workerman/Autoloader.php';
    / \2 M$ i8 j- K; I
  3. : U7 N1 @; C9 x' ]1 [
  4. $worker = new Worker();" e& Z: M) T8 S) T
  5. // 4个进程
    ( Y% V9 i, u" B7 i
  6. $worker->count = 4;
    4 z4 `+ q# m- V0 b
  7. // 每个进程启动后在当前进程新增一个Worker监听; K- V' a6 \5 Z- l1 u8 c9 _0 Z. z
  8. $worker->onWorkerStart = function($worker)' @# k3 `# g8 }- z' A+ Q+ O
  9. {  o" W2 T- _6 I0 R, a  Q3 P0 R/ l
  10.     /**
    2 G, M, C" m1 s3 c2 L) g
  11.      * 4个进程启动的时候都创建2016端口的Worker
    8 m! V1 Q' h9 S: F4 \; T/ I$ X
  12.      * 当执行到worker->listen()时会报Address already in use错误7 B; o7 F% \* Z8 x
  13.      * 如果worker->count=1则不会报错
    : ~- D! t% w$ `8 x4 v0 @
  14.      */
    3 ?/ k3 B) l4 J% ~9 R1 k
  15.     $inner_worker = new Worker('http://0.0.0.0:2016');
    * `, C4 I, U( X4 z+ x6 p
  16.     $inner_worker->onMessage = 'on_message';
    0 U/ \- X. D) [7 X0 a9 ^1 M& e
  17.     // 执行监听。这里会报Address already in use错误
    ' @  t' V% @/ K8 Y
  18.     $inner_worker->listen();5 G3 h8 O2 d, u" \
  19. };
    , `, W, A3 p: o# C  M$ h

  20. 9 e: a, a. N9 t4 J
  21. $worker->onMessage = 'on_message';( Z) l3 E5 w$ `: i) {
  22. ' R& S3 D; t! j  C
  23. function on_message($connection, $data)
    . t. I, n/ ^1 j
  24. {
    * R! z4 j0 N' ~$ N; z2 m
  25.     $connection->send("hello\n");
    # l0 F. i! E$ k) W8 R& ~! E" l9 y
  26. }
    8 Z2 r: _9 S& I( f% k

  27. : ^) Q3 K% K1 K$ A4 j
  28. // 运行worker
    ' }* {' s1 j7 ^% v1 i! E" }
  29. Worker::runAll();. V( G! g1 i. w# k
  30. 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
    * [' f! R9 g2 c1 `/ ~2 h' Y

  31. 2 T1 K2 V* k4 k# W' F2 k
  32. use Workerman\Worker;
    ' r/ |6 ?0 [0 z0 Y% B6 i# L4 x0 q
  33. require_once './Workerman/Autoloader.php';
    % w7 ]7 V0 r3 z8 {" [/ R. F) g2 E; J
  34. 9 _9 B# F( U, B* |% ^% N) W/ }
  35. $worker = new Worker('text://0.0.0.0:2015');
    1 `5 c3 H& y. x9 {* u1 q
  36. // 4个进程
    5 S8 }# q, ~0 [* c
  37. $worker->count = 4;. r( v% W5 q' c0 B* |
  38. // 每个进程启动后在当前进程新增一个Worker监听8 g& W3 K9 M  Q2 ~
  39. $worker->onWorkerStart = function($worker)+ m3 x, N" k& K9 _+ `% I+ P
  40. {/ R+ C% u$ }8 L/ d% Y% q" ?. }
  41.     $inner_worker = new Worker('http://0.0.0.0:2016');
    0 }0 u, D7 B/ N1 }5 g( L! B  b
  42.     // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
    6 g+ X7 h+ Z  e7 s( p
  43.     $inner_worker->reusePort = true;( \5 v1 ?  F8 ?' a4 r
  44.     $inner_worker->onMessage = 'on_message';3 X9 X- r0 |* d; z" G
  45.     // 执行监听。正常监听不会报错
    " e  }/ v- @4 E6 m
  46.     $inner_worker->listen();/ {6 E( {2 }; V3 F
  47. };3 s  G( h4 Q* n1 S+ n- s; C

  48. & U# m# O" W8 R" p
  49. $worker->onMessage = 'on_message';
    9 |% _$ t$ z& F! O' E

  50. 8 e: u/ j/ M& p* y
  51. function on_message($connection, $data)" ^: U7 S: \9 s0 y" l) q1 C1 B  b
  52. {
    7 r: O4 C2 U4 E/ v
  53.     $connection->send("hello\n");+ O! B; P" c  t+ M4 p/ G
  54. }
    * k7 O: Y& j& N( H4 k" M8 A& b
  55.   p, ~" m* J8 ^9 v! s1 Y  h; X5 J
  56. // 运行worker% I/ C) r7 [8 @( Y  u
  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
    " L% P9 W6 V3 h: v7 M
  2. use Workerman\Worker;
    + w; c, i% w  u+ E' E( ]6 s  e
  3. require_once './Workerman/Autoloader.php';
    , K; |9 O* c0 g% J
  4. // 初始化一个worker容器,监听1234端口2 R' N' e5 m8 H
  5. $worker = new Worker('websocket://0.0.0.0:1234');* e* i5 y$ D7 N4 Q, p8 h3 t6 X

  6. 7 X3 K3 [) d/ c) b/ V7 H  h
  7. /*+ o0 d+ L, I( Z7 I( d; g9 p
  8. * 注意这里进程数必须设置为1,否则会报端口占用错误" s: H. F4 n  M; E+ b6 n
  9. * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)% _5 O+ E% ~6 N0 y; g
  10. */0 G) A8 g$ z* s
  11. $worker->count = 1;" b& G$ U: J" a9 Q/ v. ?8 R
  12. // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
    0 F& l0 g2 D1 U" y
  13. $worker->onWorkerStart = function($worker)1 X0 q" B+ G$ B5 m! H+ l$ x
  14. {
    7 Z4 d% t- T8 ?0 O$ {* d
  15.     // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
    ' w" G( H+ D' w. v) a( V8 T
  16.     $inner_text_worker = new Worker('text://0.0.0.0:5678');
    1 G' L; b2 u! I
  17.     $inner_text_worker->onMessage = function($connection, $buffer), H6 Q% [7 V( Z0 I5 q0 [
  18.     {
    ' _/ c; _/ g4 O. Y; S- V
  19.         // $data数组格式,里面有uid,表示向那个uid的页面推送数据% a. a/ m# Z3 i5 f
  20.         $data = json_decode($buffer, true);
    * K# m% F& w/ T; \' [0 l6 @. \7 R
  21.         $uid = $data['uid'];
    ) Q' U0 S4 v4 `8 E, s4 M
  22.         // 通过workerman,向uid的页面推送数据: R& Y) s0 c' x9 n! x, J6 D$ s
  23.         $ret = sendMessageByUid($uid, $buffer);
      G; k& P. ^" T- Y  M( _3 @
  24.         // 返回推送结果8 O& T7 B7 v3 I3 R! ?9 l$ J
  25.         $connection->send($ret ? 'ok' : 'fail');
    # p  X, u" v2 E$ m3 O
  26.     };9 u# Z! ]9 v1 T8 M% C# H; ?* ~& m* h
  27.     // ## 执行监听 ##+ p/ ~8 U4 ]! f/ q! X8 I0 m, G
  28.     $inner_text_worker->listen();
    . @7 e. j% a7 D7 }' k6 S4 x
  29. };6 S  S) x# s( M9 g2 K2 }: X2 H1 _7 j
  30. // 新增加一个属性,用来保存uid到connection的映射
    % Z" @- \- I* h) M
  31. $worker->uidConnections = array();
    ; i. H% ]! c* D4 S) w; M7 j
  32. // 当有客户端发来消息时执行的回调函数* a4 q& H3 a) O0 |( n
  33. $worker->onMessage = function($connection, $data)( G1 v: c; V: m1 D" f3 z
  34. {
    / i+ ]7 G2 L4 h' `8 R
  35.     global $worker;( M) D3 b# k: s- c) H
  36.     // 判断当前客户端是否已经验证,既是否设置了uid5 K# Z. k& X) A! B$ j1 {. q! t
  37.     if(!isset($connection->uid))! F1 M! K- s( f4 h4 W
  38.     {- t* X: f  j" `
  39.        // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
    ) o5 C( ~$ ?" G8 H1 p
  40.        $connection->uid = $data;
    . Z5 R* N6 G1 b2 l4 |  i
  41.        /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
    4 O! j+ J7 ^) ^/ X
  42.         * 实现针对特定uid推送数据
    2 |! G9 p8 F5 Z' Z+ t; n8 I
  43.         */
    ! `/ e$ D& o+ F8 J, w# S8 j6 O
  44.        $worker->uidConnections[$connection->uid] = $connection;
    * x( S9 |  X; H0 d4 ]
  45.        return;  O1 Q( W  g% Q3 t2 g  J
  46.     }+ f6 Y% A" U  l1 K
  47. };* @' F, T8 b9 w: |" v+ v
  48. : R0 p/ H) a0 V0 `% A6 R8 E1 h, P" h
  49. // 当有客户端连接断开时! c, ?# i, m1 F
  50. $worker->onClose = function($connection)5 M( l4 D7 D5 O' X; b
  51. {
    4 `# \* z7 A( Q' p* d: A2 S
  52.     global $worker;- r( q6 W1 x0 e  l" ]
  53.     if(isset($connection->uid))
    2 U% C: d: i: S5 c& L8 Y4 E6 E7 q% D
  54.     {; g3 j* G' w3 ~( K7 r/ c
  55.         // 连接断开时删除映射
    " `' z7 K) Z+ P2 H  J
  56.         unset($worker->uidConnections[$connection->uid]);
    , N- y! S; @6 a' p5 q
  57.     }3 p: j' r9 T/ W
  58. };
    8 d( ~" O( t) s) k" S
  59. 3 |  J7 y6 @1 v
  60. // 向所有验证的用户推送数据6 y, |: r0 r9 A7 M0 n# c
  61. function broadcast($message)
    # K* _; u, E8 ~( Z/ y3 }3 h3 I
  62. {
    $ `0 x; T- o0 A! H! s; |3 U- e' n3 W
  63.    global $worker;
    ; ~8 ?2 M8 Y6 K, q
  64.    foreach($worker->uidConnections as $connection)4 d+ [/ R: X6 g! x3 B
  65.    {
    6 H1 r* o( d+ j( p3 E2 q- _
  66.         $connection->send($message);
    ( e% A- I! t! B& ^
  67.    }
      s( A  K) N. @# \/ g
  68. }
    ; @/ q8 o# P' l$ _5 h

  69. . C  p& l1 {$ a, c0 J& }
  70. // 针对uid推送数据; s/ o1 v. l+ a7 t7 u
  71. function sendMessageByUid($uid, $message)( O; k) r5 @& L2 C
  72. {
    ' A, ~* C; z0 G: d3 v" k
  73.     global $worker;
    ; _  e" @; D) @, I: i3 v: \
  74.     if(isset($worker->uidConnections[$uid]))
    4 Z* t2 P. f- C. ?. F
  75.     {5 b- {3 `7 X: o( |4 @( Q
  76.         $connection = $worker->uidConnections[$uid];9 D6 H8 n' ?9 Z/ R" _* I, V) w
  77.         $connection->send($message);
    ; e& W* f& Y* M4 {) a9 p
  78.         return true;
    ; d: h& p8 e+ z& }! i" F# q# z6 f
  79.     }6 B+ J& d& y. G! k
  80.     return false;
    3 C1 J2 g4 G" k: G6 v! f2 U0 V
  81. }" I; Y  o4 G4 e9 y- L3 C9 A/ d5 n

  82. 8 O2 ]$ M/ ^- G+ |5 I9 C. A
  83. // 运行所有的worker1 {8 A) o" U" [1 }( E$ E9 ~
  84. Worker::runAll();
复制代码
启动后端服务 php push.php start -d
前端接收推送的js代码
  1. var ws = new WebSocket('ws://127.0.0.1:1234');
    0 p; Y. b6 v  S8 l: Y8 M
  2. ws.onopen = function(){
    4 _' Y( ?: B+ W( n7 w) D: x
  3.     var uid = 'uid1';
    ) F" h. Q: s1 B7 P3 E) |, h
  4.     ws.send(uid);4 N7 S8 Q, @3 W
  5. };5 y. A8 [+ v  K0 p( I- S5 F( ~
  6. ws.onmessage = function(e){
    3 W# h* g7 N+ c# `  z9 v
  7.     alert(e.data);
    5 |3 }" F& `4 B! ]
  8. };
复制代码
后端推送消息的代码
  1. // 建立socket连接到内部推送端口
    & W3 u& T6 H  w- X0 P  R5 p8 P
  2. $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);% v* v) \, R* F: T# D
  3. // 推送的数据,包含uid字段,表示是给这个uid推送* w; F3 |9 f4 o3 f( c, Y- F
  4. $data = array('uid'=>'uid1', 'percent'=>'88%');5 L6 J6 v8 B( k2 ?
  5. // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
    9 R1 @9 Y4 k; u
  6. fwrite($client, json_encode($data)."\n");# x8 D  |( E" o
  7. // 读取推送结果, u0 w, o4 `9 F6 O
  8. echo fread($client, 8192);
复制代码
; F" K5 _. g& n

( G7 _. r1 v- f! \
分享到:  QQ好友和群QQ好友和群 QQ空间QQ空间 腾讯微博腾讯微博 腾讯朋友腾讯朋友
收藏收藏 分享分享 支持支持 反对反对
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

GMT+8, 2026-8-4 12:40 , Processed in 0.058797 second(s), 20 queries .

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