- 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错误。例如下面的代码是无法运行的。 - use Workerman\Worker;/ K0 e4 V9 y3 g! [2 w0 A
- require_once __DIR__ . '/Workerman/Autoloader.php';& y! Z+ {( |1 `- {
' X; i. |. W+ z; V3 e* `- $worker = new Worker();
9 V C+ Q- r7 A+ o4 P8 L4 W" I - // 4个进程+ `' U1 P# d# ~. z
- $worker->count = 4;6 p' |. R& w" @
- // 每个进程启动后在当前进程新增一个Worker监听
- {! _7 u$ i9 ?2 u - $worker->onWorkerStart = function($worker)
/ X2 B! `4 e. A5 ?$ U4 P3 _- I" g" z4 ] - {6 d2 z; _/ Y! F, F) s, N; L. p0 ~
- /**( h+ Y- m$ ]* _+ R& D, E9 j
- * 4个进程启动的时候都创建2016端口的Worker4 ^! i; ^( N. `" d4 T
- * 当执行到worker->listen()时会报Address already in use错误
! N8 d( _# N% V/ c( b7 x, i - * 如果worker->count=1则不会报错
+ Z+ {5 }: W& a. q- } - */$ x! `$ G0 w# `% J# m( S( }
- $inner_worker = new Worker('http://0.0.0.0:2016');& S! b% Z2 z. c0 Q; K _
- $inner_worker->onMessage = 'on_message';) Z1 C; |1 ^5 P: u6 K! Q; d9 P, ]
- // 执行监听。这里会报Address already in use错误
) O+ A* U5 d5 a - $inner_worker->listen();! L4 M! ]. o' R1 r2 A2 |# _
- };
& k2 _0 I- m4 l8 n - $ u" r, C1 b) m% _
- $worker->onMessage = 'on_message';/ G0 R$ M1 Q( @* j0 ^4 _( P
- " m. r+ g4 z& I& c5 w$ U, I/ O8 P
- function on_message($connection, $data)
) B9 ?$ I+ G; ^/ D7 p2 w - {6 p& D9 L/ A5 Q3 i* k/ a5 r+ c( R
- $connection->send("hello\n");% S1 G. n8 r2 N4 S* l1 [
- }
5 d; I' i$ z! o - 0 I0 i' g: ~7 y: y
- // 运行worker
& U. z. {! c- B1 `: {2 G G - Worker::runAll();
2 O! |0 x% z4 U& j+ |% y% p C - 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
" \* o4 @, ?) O& Q: @6 b7 u
3 }' c* n6 s0 k1 ^- use Workerman\Worker;
: \! o& ~+ Q$ r" N# A% o7 z - require_once './Workerman/Autoloader.php';* x& j/ W+ g$ j4 {3 j& X: H
* ^4 x; ?* ^, W; {% B- $worker = new Worker('text://0.0.0.0:2015');
7 e- C9 d0 H9 ~1 h: z - // 4个进程
! |; h3 Z: z0 o. z0 P n - $worker->count = 4;
% E/ Z$ j- s% u; e - // 每个进程启动后在当前进程新增一个Worker监听
$ p+ X. L. ~. A8 [; A$ P1 V W - $worker->onWorkerStart = function($worker)
* s h+ Y) K+ x6 h8 p& }) }8 a - {# ~" f/ E/ h2 c! p! A$ j6 c
- $inner_worker = new Worker('http://0.0.0.0:2016');2 l) H6 F! g* }) Z- q5 V
- // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
& J: P6 |0 r& W) Y; T/ y# P; ? - $inner_worker->reusePort = true;
/ l% k2 |) `3 ^ - $inner_worker->onMessage = 'on_message';( t) G$ b, J( B7 D
- // 执行监听。正常监听不会报错% ~8 u5 Q. T. O' @) R$ Z
- $inner_worker->listen();
, }5 }! E: q, B/ R# }3 C - };6 l. w* V" }+ C6 v8 ?
- 7 y! B5 {( Y( P8 W* w
- $worker->onMessage = 'on_message';
4 X( H1 R7 x* k9 V- {" p. O - ( l7 n) @" R/ G+ x" d1 C
- function on_message($connection, $data)
9 @: J6 R3 } d( m - {
3 {1 T+ S; l" m. k% ]3 z; s Q - $connection->send("hello\n");- z H& \8 H; U9 ~" m! Z9 B( |
- }3 A' A0 @; B( L) y, l# Q
- & D- }" n* ~; }& `
- // 运行worker4 a; Y2 Y! T/ W. P! t' H! k# @6 E( D
- 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 - <?php: c- R( Q# P9 {) s. |1 T# a
- use Workerman\Worker;
+ A2 \5 a8 A' J% e- \/ F! w5 A0 j - require_once './Workerman/Autoloader.php';
6 R/ ?4 x; j9 S$ f6 L - // 初始化一个worker容器,监听1234端口' n7 v& k% ]4 J
- $worker = new Worker('websocket://0.0.0.0:1234');
! c/ a0 x6 l8 @ - 0 z+ D% O9 ^; l8 p
- /*: z/ `& |# |6 g6 W
- * 注意这里进程数必须设置为1,否则会报端口占用错误
8 y; g( X( g4 x7 I1 A - * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
/ H; {9 [1 [( x/ G* x+ R0 w - */9 ~% I) v4 |* f
- $worker->count = 1;
) ^3 Y2 X! O* v5 p9 {& @ - // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
: j! C! Y2 }$ y7 O( k - $worker->onWorkerStart = function($worker)
$ s; t+ v6 c: m! [ - {
+ w- Z( C8 s0 j& b( i5 X: h$ N9 ] - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
8 ]+ M6 t% ?$ i; g - $inner_text_worker = new Worker('text://0.0.0.0:5678');
# F4 i: \3 x, h/ r) o - $inner_text_worker->onMessage = function($connection, $buffer)0 w V5 D, K: F4 b5 {, A
- {
% ]/ o2 N0 h2 u' k2 r - // $data数组格式,里面有uid,表示向那个uid的页面推送数据
0 E/ e+ }% b( @9 n% N S. y$ ? - $data = json_decode($buffer, true);
" h) h3 N7 }5 e/ x& h, g - $uid = $data['uid'];$ D* K9 ]8 s! w4 `# X$ v7 F4 B
- // 通过workerman,向uid的页面推送数据" p+ R( ?) [4 M
- $ret = sendMessageByUid($uid, $buffer);
! V2 L! g/ X' T* }6 H/ S1 \0 s - // 返回推送结果' D1 x1 `6 L! x8 o T, E
- $connection->send($ret ? 'ok' : 'fail');
6 a8 w ~" h/ F N. P1 z - };
- Q0 z" j7 Z9 _9 G - // ## 执行监听 ##
$ J+ O7 D. V0 R8 L% h: c/ A - $inner_text_worker->listen();
Q8 E. s6 Y1 k- K8 m* W! F: O7 B0 ~ { - };, X7 a8 L n4 R, d8 c
- // 新增加一个属性,用来保存uid到connection的映射: E8 r* z" u$ D
- $worker->uidConnections = array();
" Q+ A$ h9 p9 t. g8 P5 F( n$ E - // 当有客户端发来消息时执行的回调函数
8 a; o, w" _5 n+ r0 ^: ?4 R - $worker->onMessage = function($connection, $data)
/ s2 N2 L5 W0 l" V( i+ b1 V - {) T h7 l! n' u; b1 B7 ]$ g+ `
- global $worker;
+ Q# t% Z6 j- b' V6 C - // 判断当前客户端是否已经验证,既是否设置了uid! ~9 }9 g M! @* X
- if(!isset($connection->uid))
, S- k9 u" R5 t1 t7 c5 C( R - {
! u) O5 b O* Y G - // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证); i) |. I$ p9 B7 p* P$ i+ `
- $connection->uid = $data;
/ q, P6 V( q8 M0 x# ]" t f( ^ - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
2 _9 m, L3 y W- z0 u5 ^6 }% u - * 实现针对特定uid推送数据
0 z- u' T3 p4 \8 J' a( m - */
+ w& }) M# o; D& e - $worker->uidConnections[$connection->uid] = $connection;7 P/ y0 g2 z& F+ v* |. Y
- return;$ `7 i7 G# u( }( L; y9 K7 q
- }
, A" y: a* k2 m. ~, x( D - };
: t9 }$ z0 v1 N( u. v8 g - 2 ^( f, R9 f* J: Y4 q
- // 当有客户端连接断开时8 @; \8 L* x) p" T! j: A O
- $worker->onClose = function($connection)
' O8 T2 h o+ b. C: d8 x - {
3 A; p* o( L4 i1 {& X - global $worker;
9 i& S5 S! Q' Q- P2 e/ @1 y - if(isset($connection->uid))
7 k4 ^/ I8 O& R4 e - {
. q- t" K5 d) `) a5 q - // 连接断开时删除映射
) ?+ A+ ^8 _% {: M' n - unset($worker->uidConnections[$connection->uid]);
) l, y7 Z6 \# T; V/ } - }- D$ W' A, _3 n4 {
- };3 e. y" }% H# J [: Y' S* ]- b
1 |7 S" ~ I( @7 g6 y- // 向所有验证的用户推送数据
: O5 j2 S+ h9 U- w; M* `& y - function broadcast($message)
) t# U; I; @2 f5 t: ]6 E - {! v# \$ x, P' o! y& N0 Z9 S6 X
- global $worker;
# Z" ]) k S: M( K, q# L - foreach($worker->uidConnections as $connection)
o% X) ?' a7 W# v& t1 b# u - {5 D: F6 }. x8 O7 k0 O; u
- $connection->send($message);* r8 O; R; _% p6 `- b- L& o$ c4 x
- }- c$ Y3 h8 {8 o2 @% d- j: E( d
- }
& x) f' r+ g$ S2 C, T- i3 ^* s% l
; [2 b# F+ g# b! A- // 针对uid推送数据
9 z1 U2 g: s, m8 h9 Q - function sendMessageByUid($uid, $message)8 d; c# \+ K9 {" `
- {0 A" u/ b1 V B$ s* b' {
- global $worker;
$ K& q3 I* T1 y- m% Y - if(isset($worker->uidConnections[$uid]))
+ \" U3 a$ O( m+ [! s4 C$ f3 b* i - {
1 g# n. z0 D: z6 x& i2 } - $connection = $worker->uidConnections[$uid];$ W$ z" L, D \2 c
- $connection->send($message);# O2 q0 {; }; {7 b i: `0 h/ E
- return true;$ K8 c4 S e Q8 A
- }. O7 L% ~* k" }, T1 E" V- Z# k- I/ Y
- return false;
. K7 B# H8 s) p/ E$ k6 Z - } B) N" s+ ^4 O% Z- i+ h
. T! z6 F2 l! e0 ^6 E5 _1 I Z- // 运行所有的worker: Z4 H8 q3 f/ V6 |
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');
% @# g. h/ z: j( Y3 j+ K - ws.onopen = function(){
2 _) i Q3 A' z& w- Q2 y - var uid = 'uid1';
8 p \' c3 s5 J9 V& j - ws.send(uid);
* Y; v+ l4 r- \' V7 Y3 I - };6 x X& Q0 i' b& D) H; E( R
- ws.onmessage = function(e){+ I+ ?/ H+ w, e: B8 n7 i
- alert(e.data);3 d: B W; C- N# r r2 w! @0 X8 j/ D
- };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口, X% J2 D5 g% _- o+ x- L$ l
- $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
6 _6 D/ q/ ?# S7 _ - // 推送的数据,包含uid字段,表示是给这个uid推送
) ]% T$ q3 }2 F7 T* V; ^ - $data = array('uid'=>'uid1', 'percent'=>'88%');* ~7 N3 e0 X7 A9 h" _; R% }
- // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
! [" l! r5 L& |8 t6 Q: { - fwrite($client, json_encode($data)."\n");
0 a! t/ X" s8 Z/ r/ Q: A3 v" } - // 读取推送结果
% @; N- Z9 R: |. W - echo fread($client, 8192);
复制代码
# |) y6 m. o, G. F9 ?. ?4 O7 J5 ]! C) H
|