- 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;
- s+ x6 d; F. u1 f( @8 b - require_once __DIR__ . '/Workerman/Autoloader.php';+ f" Y3 [6 v* r* R. D4 R
% P3 p! {: Z5 N) `, ~4 d- $worker = new Worker();
$ @. i, u( ?9 O! E: T2 L) h - // 4个进程
8 X# U9 U1 Q6 l8 {& i' N - $worker->count = 4;
/ E/ G" j; U! C2 q2 ] - // 每个进程启动后在当前进程新增一个Worker监听" N6 g8 [) j$ w& y% Q4 n9 l
- $worker->onWorkerStart = function($worker)/ g+ ^! m8 @0 w: z9 {
- {/ U; @* q, k# A2 a% {# Y7 E' s7 x
- /**
3 A0 W( J. V0 }- d5 V& g - * 4个进程启动的时候都创建2016端口的Worker/ ~. X0 f2 S. V# S# @
- * 当执行到worker->listen()时会报Address already in use错误% ^8 z# h' [- H
- * 如果worker->count=1则不会报错) J, Y' Q5 E8 g: ]
- */
8 r0 |6 `' s) g% K: N% \ - $inner_worker = new Worker('http://0.0.0.0:2016');* R9 ^' E) u6 |. V! |
- $inner_worker->onMessage = 'on_message';
8 z4 ]' X- `- X7 s% l. A; F s - // 执行监听。这里会报Address already in use错误9 D- D3 s+ J8 [' j
- $inner_worker->listen();
5 ~ f" {$ ]$ s - };
! ~8 s2 n6 U( i% _
$ J6 e' ]4 q# O% T; W: P" E9 v1 M) b- $worker->onMessage = 'on_message'; Q$ X/ u6 v; `& ?% G2 Y$ s
. g' X7 l$ T* `/ b D% p& R" \- function on_message($connection, $data)2 s: _' U2 J; m
- { _& N: z/ x4 l# O$ u
- $connection->send("hello\n");
9 Y" B6 g V6 Y - }
0 d. G7 Q2 k( t v0 c+ Z8 t" o" c
5 R6 j6 I* n+ S7 j: J0 s4 ^& s- // 运行worker
" T% V; D; `5 ? - Worker::runAll();) M7 Y/ j" q2 d) Q' z* R) K
- 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
0 f) l" ]6 ?& i! v% d$ y, F Y - 6 J3 `' r, c4 x2 ^1 y1 h! _
- use Workerman\Worker;) ~- w7 f8 P+ r8 r4 c
- require_once './Workerman/Autoloader.php';
6 X& _* H* d4 Y/ {( l# ~1 A* I
3 o3 p$ v5 O$ E- $worker = new Worker('text://0.0.0.0:2015');
4 X+ k1 f0 T1 j8 ]* _" t - // 4个进程& g! n4 L9 U& F n2 i1 l: v- v
- $worker->count = 4;" V: X1 P& l2 A
- // 每个进程启动后在当前进程新增一个Worker监听
3 ]- M! D( {* D- {6 ^- g! w; i - $worker->onWorkerStart = function($worker)
1 s6 c5 p( B% J L) T |: v - {
5 M' w% b/ T" x - $inner_worker = new Worker('http://0.0.0.0:2016');. b3 T/ h7 u8 o2 x: J5 q
- // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)! X% G. J8 [, e) ~' |
- $inner_worker->reusePort = true;9 @. G4 y) N' O, Q
- $inner_worker->onMessage = 'on_message';
; A) _: b9 j$ N# _; f/ B! H - // 执行监听。正常监听不会报错& a# N7 V- ~* {- F# E$ k" |
- $inner_worker->listen();3 L! Z* k, P: |8 Z+ X W& T
- };
0 B" d$ A& _. N
5 H5 z4 V0 [: s" \- $worker->onMessage = 'on_message';
, q: m4 i$ r6 \- Q/ T# X - . |1 L. J) G; w7 y% B) Z+ a) `
- function on_message($connection, $data)
3 J0 z! `3 |+ N) B - {9 J' k/ R: H) Q- T
- $connection->send("hello\n");' ^# J1 [' c6 Q
- }* k+ U" X: X) F7 p
- * a& o+ v* n, S8 K O$ P H) P: n
- // 运行worker
/ n6 K- A* L6 ?9 R - 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* {: P4 k2 I1 o$ k* E, a* {
- use Workerman\Worker;
+ d! y8 G; w$ H. j N0 J4 M O - require_once './Workerman/Autoloader.php'; J8 O& a8 [- a$ \" q1 {* m; i' \" K
- // 初始化一个worker容器,监听1234端口8 p4 p: y7 |0 Q7 {$ I+ l4 b
- $worker = new Worker('websocket://0.0.0.0:1234');
1 Q. e* n$ r! Y( H. f - - C5 f8 l. ~/ t d* L5 C& x
- /*# |0 R7 j' b& g6 N3 g
- * 注意这里进程数必须设置为1,否则会报端口占用错误: r P1 \+ X5 D4 F6 u
- * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
9 p5 M- L* {0 ~+ q9 V4 l+ Y - */. V# ^ N* r7 {, w3 y5 a
- $worker->count = 1;1 J: J7 ^9 N" X2 B+ M# N/ M
- // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
' i9 i3 Z; F; q6 C - $worker->onWorkerStart = function($worker)
9 ]5 |4 Q3 o \1 K. _ a - {, c; j3 B; y6 C: D
- // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
2 @; H# S8 V9 P9 \. Z5 t - $inner_text_worker = new Worker('text://0.0.0.0:5678');
% q/ X4 m* W& \) h# g - $inner_text_worker->onMessage = function($connection, $buffer)
" w2 U1 \- b$ |4 j/ c+ U - {+ D+ ^3 P& o {1 }
- // $data数组格式,里面有uid,表示向那个uid的页面推送数据
# e/ N9 P/ c9 C0 s' t V: _' a3 J - $data = json_decode($buffer, true);1 t9 Y3 C! Y5 I! X
- $uid = $data['uid']; |8 w6 u: o9 v @7 a2 C
- // 通过workerman,向uid的页面推送数据2 a7 L+ |$ G4 E. r& A3 D
- $ret = sendMessageByUid($uid, $buffer);
& V/ V& x' m9 x1 |1 S4 v( _ - // 返回推送结果& C9 Y2 b" ^7 C4 F! ]; e' U2 w9 G
- $connection->send($ret ? 'ok' : 'fail');
2 b$ s7 @- A7 X8 {3 q- B; J; P - };
! i) l* `# E! z" Q - // ## 执行监听 ##
3 ~. ]- ?6 A. J7 y$ y8 q - $inner_text_worker->listen();
( S1 A8 W3 |3 E1 j! [" B - };
5 _1 E( C$ I; x - // 新增加一个属性,用来保存uid到connection的映射! I; A/ J+ z3 F# X6 o* U
- $worker->uidConnections = array();
5 H9 p# \9 z& r+ I - // 当有客户端发来消息时执行的回调函数
9 w6 T" O7 U" ?" | - $worker->onMessage = function($connection, $data)0 F) k$ L; @# N6 K# K
- {% L( P8 g; _2 Q/ C0 D
- global $worker;
. Y+ k$ }* l6 s! c7 |, {5 H - // 判断当前客户端是否已经验证,既是否设置了uid4 `1 b, l1 l/ e" O6 v2 a
- if(!isset($connection->uid))
I3 A# R, B% U' X( J7 ? - {: S$ Q5 J0 e9 P
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
/ `2 r1 ?. I. @. E m% u - $connection->uid = $data;$ E) [" y" P; A
- /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,) n) B4 P8 i2 A, I; @3 e$ d
- * 实现针对特定uid推送数据
. v) h! b: \) \! P6 U( X - */
% P# M; j0 r7 Z0 v - $worker->uidConnections[$connection->uid] = $connection;5 p9 s# S* o7 o% P& M' k
- return;. l9 E6 K! j+ Y2 ^: S0 j/ U: j( p
- }1 l4 O& M6 F" T0 z4 f; U
- }; D5 d4 R/ h3 J( e& P$ V ]
- 2 A" E D; v7 @5 m) y$ i, ^, }
- // 当有客户端连接断开时
6 C5 c; V% L9 p9 s! ~( ?7 M3 Y9 W/ \ - $worker->onClose = function($connection)2 h" h8 k6 b8 j% W# m# M: m( z
- {: M+ m. f$ W4 }3 N
- global $worker;, [/ x+ V( ? L+ b. O5 k- Q
- if(isset($connection->uid))+ B8 b/ ]1 @5 q2 D
- { I, J" X+ ~4 f" e4 c
- // 连接断开时删除映射, x: A6 }4 S v* s: v4 l
- unset($worker->uidConnections[$connection->uid]);' N' z2 y) ^: _2 |$ j8 t
- }
; v+ [6 N$ e8 a. U& P+ |. G( k1 z- B - }; D3 {9 S& k) b- P. ?1 ~
- 6 J: Z( e5 { K4 C- E8 x2 d( F1 Y* O
- // 向所有验证的用户推送数据( P, u, I# H/ z* N' Q n: T
- function broadcast($message)
2 A9 K" ^7 X, [; h2 B0 o - {
( n0 v# X9 S0 `) \+ A o - global $worker;
$ V, u) }" @/ f9 U& J/ _0 d - foreach($worker->uidConnections as $connection)
* E3 L6 r6 U1 z4 \2 {. D - {
% K/ V( L# ~/ A6 g* s$ ^' o - $connection->send($message);
6 ~) i9 W+ O' U4 a" K) g: W) z - }
8 @6 M, k* k; @ - }& k2 Y& P1 @1 I: |: e
# I" H2 G+ l' w- // 针对uid推送数据' z8 H- z2 r% N7 G
- function sendMessageByUid($uid, $message)# U; E7 @8 E& Y6 A4 ~
- {9 f7 E1 Y6 A4 p8 U8 Q2 s
- global $worker;4 i: {9 `& X! c/ [) r: p* X, J
- if(isset($worker->uidConnections[$uid]))
$ q$ C, E' K9 d* Q* J - {
" n' [+ ^/ T; O8 [6 e* X- ^4 i( F1 p - $connection = $worker->uidConnections[$uid];3 i7 x6 v* P9 G) i0 d
- $connection->send($message);1 C# \; M4 y3 a! r0 C
- return true;
, d+ ^( o9 v1 o2 M7 x - }
4 _! T- \8 B- v4 L J% A% b4 Z - return false;# G6 G; `% }. y h `" @1 ^
- }6 l# R% Z7 M1 d
- 0 X- v: `5 d3 ~% ~( a* q1 _
- // 运行所有的worker5 g% x% [6 [/ S& N" p" l1 P5 W" l
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');
0 v' ]+ C: e" ~1 c6 T; F# f - ws.onopen = function(){- Q+ y& q( t. B% _
- var uid = 'uid1';
) B* k3 A4 A5 S9 g0 _ - ws.send(uid);
: o' k5 e# G8 O' u/ j; i - };6 Q0 ]7 B0 g' w3 p3 s/ n& d/ J8 H
- ws.onmessage = function(e){0 c' @% Z1 u; f, R- U6 N* Q5 {
- alert(e.data);' p. m/ z, B" S+ a" K( L
- };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
; g, q" R* y6 s/ ]; H - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
' I! M" P1 g9 a4 z2 q/ O( r - // 推送的数据,包含uid字段,表示是给这个uid推送- w/ p" _( Y3 _2 K
- $data = array('uid'=>'uid1', 'percent'=>'88%');
# q$ A4 v$ Y& C' E5 r& @ - // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符( Y( l" R# F1 ?- H) ~6 x
- fwrite($client, json_encode($data)."\n");2 z- A/ t- B" S9 B9 k: H! L
- // 读取推送结果& ~6 y+ P& q/ a O7 Z( ~2 O
- echo fread($client, 8192);
复制代码
. f3 Z" U( m: y' [
: d3 K( h( ?% ^9 A |