- 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$ v7 K* ?, Q7 v! N0 l - require_once __DIR__ . '/Workerman/Autoloader.php';
: [$ N: P& o; N1 H
; K. f: r7 t& u* g0 F- $worker = new Worker();
8 Q9 J h- k9 F5 U - // 4个进程# I$ \1 n* s. { {- r
- $worker->count = 4;
6 x& A/ N. E6 N/ R! P2 R - // 每个进程启动后在当前进程新增一个Worker监听& W7 Y7 i8 a& k; p/ { d
- $worker->onWorkerStart = function($worker)* ]: g' a2 W1 ]8 I* b" e4 Y
- {
4 i* y* r9 K8 m6 i6 r - /**( S3 j f1 m, d& t' G$ N( W: Q
- * 4个进程启动的时候都创建2016端口的Worker
- w* E- G+ P3 _/ B - * 当执行到worker->listen()时会报Address already in use错误
) K; m5 b' c4 m3 y; L$ \* o, p1 F - * 如果worker->count=1则不会报错: n6 o2 V+ ?& r' x3 y4 {
- */, y( r3 ]; h6 c
- $inner_worker = new Worker('http://0.0.0.0:2016');
+ l1 i2 n8 z' l6 I& ]; \% E - $inner_worker->onMessage = 'on_message';
9 I. f( z' }3 C2 a( @/ X& w; W! F - // 执行监听。这里会报Address already in use错误
7 m. b$ B+ k& u1 V - $inner_worker->listen();8 Q# {# d4 [" E/ X& F
- }; Z% W) d7 B, ^" k
- & e8 K7 U1 E1 s! l6 o9 G" o
- $worker->onMessage = 'on_message';
R6 d1 F( s: Q5 f$ _' c1 [: b% s" m - . X5 D' X2 m3 _8 r3 S8 i* [) K
- function on_message($connection, $data)
6 i. W7 s3 O. ~5 H - {
m1 p) K4 K- }5 I+ ^ - $connection->send("hello\n");5 O4 @" A* |9 L+ O. A( y2 o3 y
- }2 `3 M a" m/ n" x" k
- 9 t& d0 H% Z/ @0 @ Z
- // 运行worker: k' D* Y( e' N3 ^/ n, x9 Y
- Worker::runAll();1 b! e* c+ ]$ D0 o
- 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:' b7 I+ |) @% a# ?
- ( o2 }( C) g0 W+ v p& v' u
- use Workerman\Worker;! L8 P+ ~2 H, O: x2 K& a
- require_once './Workerman/Autoloader.php';: V3 W% K5 |6 e. d% W, n$ ^
o3 _0 s `: \+ Z/ h2 l& _- $worker = new Worker('text://0.0.0.0:2015');
4 t1 _! b7 E7 z# C K$ X - // 4个进程% ?' D4 ?* R$ y9 I9 Z$ z5 [" z6 w* ?
- $worker->count = 4;
4 ^0 u' N+ D. _ R# j - // 每个进程启动后在当前进程新增一个Worker监听0 V: w, x" f4 Y
- $worker->onWorkerStart = function($worker)
& b; Z3 c1 z } L6 z, a4 u - {
9 w4 z- v5 F& @+ W# R# m - $inner_worker = new Worker('http://0.0.0.0:2016');
* C; u9 p2 G- p" K - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)% ~9 y+ A3 H+ J5 P* D* c
- $inner_worker->reusePort = true;
7 f& a x6 a2 c9 b; g) [" f - $inner_worker->onMessage = 'on_message';
0 P+ R9 L* W: F - // 执行监听。正常监听不会报错* G V( R1 R5 X- d6 M' Y) y; J$ r
- $inner_worker->listen();/ c: s; {* v% ~- D0 c: e
- };7 u/ y" w6 Z- b6 E0 ?: ]
- % Q$ f' T% W: B0 i* P* z) R/ s
- $worker->onMessage = 'on_message';
7 j/ ?# W# o; t - - R6 d5 \4 j, W
- function on_message($connection, $data), c% _) G* G2 b- q, h
- {& U% X% ]$ }( a @1 O
- $connection->send("hello\n");: y8 k& u7 `) B g2 J2 ~4 K$ D
- }' Q1 I0 X! a0 z2 G, j( ~
- : I3 y# a- {# b) `5 l
- // 运行worker6 I9 ^ c* f. x0 P
- 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/ |6 u& A: g- E+ U
- use Workerman\Worker;
/ g% i: q; R9 S: o& b; C- K - require_once './Workerman/Autoloader.php';
- J8 r, n/ n' j. w/ s$ Z - // 初始化一个worker容器,监听1234端口/ B" I' o5 Y$ k- p1 J
- $worker = new Worker('websocket://0.0.0.0:1234');
3 I; l" M f# r9 A3 Q2 U$ [' r( N/ a
N5 k. x l! n/ O+ ^) i- /*
7 T7 \5 G9 J+ }% I1 R+ m - * 注意这里进程数必须设置为1,否则会报端口占用错误
0 Y2 o4 \% D; w. [5 Q# i( z) r1 d3 b - * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)( `. h6 a0 z J( Q. _9 Z& Y' ]
- */
' v- b1 ]3 d5 m& _" w% A - $worker->count = 1;( t- {' d2 g, S5 o3 u
- // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
- q! ?; s7 \4 a$ A+ | - $worker->onWorkerStart = function($worker)) P/ M: X# x( Z
- {
3 r# W C& v! t1 J* y - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
. A2 f6 z) N8 f6 w0 j - $inner_text_worker = new Worker('text://0.0.0.0:5678');
3 K" C, m. P j - $inner_text_worker->onMessage = function($connection, $buffer)
2 _3 u& T7 L; b9 [3 z - {+ o1 r) o$ A, L
- // $data数组格式,里面有uid,表示向那个uid的页面推送数据
/ ^% X# n7 J1 Y6 _ - $data = json_decode($buffer, true);+ X0 B7 h: k. r: b
- $uid = $data['uid'];
2 p; B7 a4 N& o, h - // 通过workerman,向uid的页面推送数据8 w9 E+ m2 h# X6 e
- $ret = sendMessageByUid($uid, $buffer);8 |) p! h+ r5 D
- // 返回推送结果
& Q7 Y; ~1 L9 m' d - $connection->send($ret ? 'ok' : 'fail');4 Z; x. y& j6 u: Z$ Q p
- };
) J$ \4 a, }3 v0 ^7 y - // ## 执行监听 ##' y9 A3 J/ A6 s$ p& \- G" P: p
- $inner_text_worker->listen();0 ? ]) |7 N. h* r7 ]* b$ g0 P- I5 T
- };, [3 o+ H: K! n& A2 m2 u* n
- // 新增加一个属性,用来保存uid到connection的映射
+ h; {: a$ L7 t- S. `8 d: D - $worker->uidConnections = array();
5 X- e8 D. |* C! m+ p/ g - // 当有客户端发来消息时执行的回调函数
5 C s) X2 m7 q - $worker->onMessage = function($connection, $data)
+ h4 l2 k/ H3 K& w( A - {' S& R J1 J: V
- global $worker;% {4 ~$ r. k/ p( g/ N: ?
- // 判断当前客户端是否已经验证,既是否设置了uid
9 F1 G8 @" V C! \ - if(!isset($connection->uid))
6 \; }" v1 e$ F0 j$ |* t/ y9 P7 T - {
$ ?& U" L. v4 V( A6 E - // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
n+ y+ E- H! I% `8 S - $connection->uid = $data;- `: B* O g R& U3 y
- /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
5 V* W4 V* ~6 c4 i9 c( _* l - * 实现针对特定uid推送数据7 w% }" C+ Y1 s+ \: U
- */
+ N5 U, e7 u$ B. x8 ^2 n - $worker->uidConnections[$connection->uid] = $connection;
8 Q4 Z- F6 Q8 e+ X - return;2 o& j; l$ E! H$ B- `1 p
- }
4 {$ h. V2 g6 n* { - };
4 o2 i) |- U; A/ @" K - 1 h: h0 N" _& T: n! k
- // 当有客户端连接断开时& z( H8 K! d! Q9 I3 L1 d& }# W
- $worker->onClose = function($connection)& M3 L/ I' O) Z5 |1 H4 g3 }) i' j
- {& E5 f! s. k3 [( |8 s
- global $worker;4 }* k9 @2 V# T
- if(isset($connection->uid))
. E, ?8 H- h0 m9 n7 h$ { - {1 M+ y. u4 \" y7 y! S
- // 连接断开时删除映射, c3 ?- y1 Z" a2 B" l6 L
- unset($worker->uidConnections[$connection->uid]);
* F8 `7 y9 m g* _. C- Z - }6 ]) Z! s1 X; }% Z9 f7 c3 G
- };
" _' k* l5 M/ s$ }8 ~& { - . n* L2 X5 b/ |5 G. U' m
- // 向所有验证的用户推送数据. C9 I: S6 C/ E. d" ~
- function broadcast($message)& X. v6 O% S1 J9 H2 m% N; d1 L6 N' K
- {
! k2 H' k @# J* t2 @6 h; m+ U. z - global $worker;& n7 o) p( a8 A( H1 v5 D
- foreach($worker->uidConnections as $connection)) ?; q2 ^& b& E& Z
- {4 h. ~6 Q( A/ b6 N4 E% D+ K3 {8 x
- $connection->send($message);
0 y, z* u: X0 [4 }( s n - }* L6 {2 Y: a* L
- }
+ L8 q( U% U; ^% Q
% ?9 a& w9 a3 X$ B- // 针对uid推送数据
' Y; G3 y8 A) ^) h - function sendMessageByUid($uid, $message)" `2 x4 }; k; x( F9 O% E- O8 D
- {
) [' i% m+ P# n/ s/ L* h - global $worker;5 f$ Q9 D! n6 G, \$ x' G
- if(isset($worker->uidConnections[$uid])); q. Y2 i5 Y7 N: k
- {
0 a8 ]8 s/ {4 T, q - $connection = $worker->uidConnections[$uid];
) M% O3 O9 A1 f% N d2 l9 G - $connection->send($message);/ n4 K7 l1 Z' {6 @; H W
- return true;
% i8 r( T& f' n* Z, M& W+ L - }
# d1 n& V' ?& N0 J - return false;
, e/ u0 R! b( e/ r3 P - }
0 G+ R& m% n! x& X
( W0 @$ l+ n5 T# X x- // 运行所有的worker b. C/ I+ T8 k4 K1 a; f: }
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');( x5 A" P& M8 k
- ws.onopen = function(){; q& Y( f4 a C: W( f
- var uid = 'uid1'; O+ H- O+ l# g: \# P
- ws.send(uid);& C6 K0 R# T. Z% v1 V& a
- };1 \+ Y% ` ?7 V0 J/ |9 A
- ws.onmessage = function(e){
1 Z# a! x; u3 ~/ u7 ] - alert(e.data);
' w! x; f, o: o$ G1 R# ~ - };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
+ z2 z A1 ?, g6 x. M5 E - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
e) X$ l4 O- Z' _: m - // 推送的数据,包含uid字段,表示是给这个uid推送
6 \" y5 x! X5 P. O! L - $data = array('uid'=>'uid1', 'percent'=>'88%');* n! J. }1 O, s5 x5 p% b: p) f
- // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
9 A |9 D% O+ f/ N& c3 D& \4 m - fwrite($client, json_encode($data)."\n");! p/ u) U# o; K3 _% s: X
- // 读取推送结果0 V+ m' }) S8 [/ F' H. n& B
- echo fread($client, 8192);
复制代码 ; }3 h, _( V0 Q
! Y3 W. _8 ^2 q# I) M( h |