- 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;, M; [7 e) i% f' a9 s# x: m* [/ N
- require_once __DIR__ . '/Workerman/Autoloader.php';
) x- [1 j( ]3 z/ S5 Y& J
7 |( w& n- R' O- $worker = new Worker();% b% {' A0 u. l
- // 4个进程! F& l" D5 P( V" X1 U Y8 O" A
- $worker->count = 4;
1 d$ I! I1 K4 T5 a$ V" F - // 每个进程启动后在当前进程新增一个Worker监听# c3 S% ]1 y; @- E
- $worker->onWorkerStart = function($worker)
C' A/ J' `& S+ F - {
w1 _* p; p" a: U - /**
, l/ t: c3 r6 U0 C - * 4个进程启动的时候都创建2016端口的Worker
4 r6 R+ @1 ^& c9 } - * 当执行到worker->listen()时会报Address already in use错误
: C4 [8 V ]- `" V% Q - * 如果worker->count=1则不会报错9 t* c# Q: t4 K
- */
h8 f' e. Q7 ^ - $inner_worker = new Worker('http://0.0.0.0:2016');
6 g8 q2 h) I6 s: {8 M) {8 |: N - $inner_worker->onMessage = 'on_message';; N4 J' J$ E) c4 t1 Y
- // 执行监听。这里会报Address already in use错误; B/ f( h* R8 Q+ K# i
- $inner_worker->listen();
& p6 d' m6 g! q- n+ P, f( @ - };3 ]2 c" ~4 k- l, O5 b% d- S
- 0 n' e; A2 F$ _5 g0 j
- $worker->onMessage = 'on_message';
% ^% L( ]$ \4 n/ F6 {: | ?4 @ - ( b+ {: K( J5 U5 K: m9 x7 ^# d( }
- function on_message($connection, $data)
' L! s4 \2 ~7 L, A - {% B, {2 M, p# X2 P1 K! V
- $connection->send("hello\n");
# o3 G0 E/ P+ ~* X" m - }
9 f, d2 p8 t. A5 q/ ~
, i) Q: n; M* I- // 运行worker/ ^. o. ]8 @. V
- Worker::runAll();1 ~2 M" h6 \# b$ @9 k% H b4 n
- 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:1 p+ ?. a c# d/ N' K
- 1 f, s# B0 p' y( L
- use Workerman\Worker;
" K. H5 J2 H5 C8 d, k; G' H5 V9 u% t - require_once './Workerman/Autoloader.php';
1 D* S$ J- q8 b# w0 O$ s# o0 s
" K, ~ x: M9 y+ p$ H) g( ]- $worker = new Worker('text://0.0.0.0:2015');: `) L& \* _! Z' J v* P. [5 g6 Z
- // 4个进程; c% @' c7 v' i d4 _( [4 l
- $worker->count = 4;+ P& H8 \/ l6 x! L
- // 每个进程启动后在当前进程新增一个Worker监听" K O! a5 O7 {
- $worker->onWorkerStart = function($worker)
! _# i7 y: w: X$ S$ c; [; b: Z - {0 v: }5 @: C2 o3 h/ Q4 Z7 f
- $inner_worker = new Worker('http://0.0.0.0:2016'); m2 G# U$ {4 V) v5 b
- // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)5 I, r+ u& D+ f8 O9 }6 F$ i
- $inner_worker->reusePort = true;$ \5 J+ b6 `% R8 C6 W( y% v
- $inner_worker->onMessage = 'on_message';) }/ K: d: ~- P
- // 执行监听。正常监听不会报错
; v& h0 Y, N# o9 ~( V8 ` - $inner_worker->listen();
, q. r, `8 F4 u5 B/ u0 V - };
3 T) M' I0 m, |; |, n5 A, Z7 H+ ^
4 F/ g6 h6 O4 ]- $worker->onMessage = 'on_message';
1 e/ [& z9 u( L - 8 E- Y# d( e1 ?8 i
- function on_message($connection, $data)
^7 W% T/ g' l - {
' k; g. f2 }# U' q - $connection->send("hello\n");
/ p) Y7 m$ }4 M. W3 g - }
& e# Z9 F6 U) t7 r* j% g* j5 K - ! b9 p8 H. E; B b% ?
- // 运行worker S+ v7 Q) J/ ^* O$ _
- 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
2 @. Y& V ^, C" x& E6 |0 L - use Workerman\Worker;
- o# N0 J6 M! n0 X5 n - require_once './Workerman/Autoloader.php';# v& S' S# s6 E: h& ]1 r& @0 c
- // 初始化一个worker容器,监听1234端口+ h+ o% r2 B6 O9 S( h
- $worker = new Worker('websocket://0.0.0.0:1234');
2 g3 U3 P6 J: ~
: W; H9 B f. c( \ e- /*0 D5 H `5 H' v! B+ T
- * 注意这里进程数必须设置为1,否则会报端口占用错误
- D% U, n7 W8 E# T2 N8 \( D4 ` - * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true), l, m0 t; a- C7 R0 @, j/ q
- */
s3 d3 R. y# y' P# N" J - $worker->count = 1;: N- ?- `+ E. Y
- // worker进程启动后创建一个text Worker以便打开一个内部通讯端口. x% u- u3 n, u$ K
- $worker->onWorkerStart = function($worker)! I E9 B3 k4 \. g7 y$ U8 M4 C
- {
# Z8 g2 `! G2 n1 s - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符7 ~( Z) N" T6 b+ _
- $inner_text_worker = new Worker('text://0.0.0.0:5678'); M T8 s2 u; s0 m# R
- $inner_text_worker->onMessage = function($connection, $buffer)) O( ^( ^; M+ R6 ^
- {6 ^5 F' e, o$ N: ~/ X& x+ D2 u& t
- // $data数组格式,里面有uid,表示向那个uid的页面推送数据
% X& @5 V6 U2 W! y - $data = json_decode($buffer, true);) H$ D- M' R6 E7 O$ E* X- ?: J6 B
- $uid = $data['uid'];
2 E6 R4 M9 W5 L' c6 O+ k1 s" g - // 通过workerman,向uid的页面推送数据- v* S3 o( c1 A$ x
- $ret = sendMessageByUid($uid, $buffer);
- }$ R, S9 b& n. F5 ~; L2 z- @ - // 返回推送结果
5 U, Z6 x4 i4 l# D: N - $connection->send($ret ? 'ok' : 'fail');! F O! N/ {5 p" W; {$ z* N" q
- };$ Y7 ~5 Y% p6 d% j
- // ## 执行监听 ##
' {. ^& {' q! D n - $inner_text_worker->listen();/ {4 O/ I9 z5 D- [; l/ j
- };; U/ ~! S$ w8 a2 b: z
- // 新增加一个属性,用来保存uid到connection的映射
7 d" x, A& V" [" |& A. L* M) g - $worker->uidConnections = array();# v* }4 q3 L- v9 S- j
- // 当有客户端发来消息时执行的回调函数0 t \7 ~* r# }5 K
- $worker->onMessage = function($connection, $data)
4 h: J& G. l# C4 W0 [3 i2 t' D' Z - {+ V3 l1 m6 H1 Q$ ^- l( r0 k8 U6 I3 n
- global $worker;
( O- I8 V1 Z; L; K - // 判断当前客户端是否已经验证,既是否设置了uid
( l+ [1 W' K* Y/ f# A, E - if(!isset($connection->uid))2 V- t2 h" O; t- U' t1 t/ c, _+ N) \
- {9 L/ f# v* k6 `- U& l0 `% \+ L
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)2 s. s6 f. s, ?
- $connection->uid = $data;
$ Y* q# c/ q& f! I- }7 S2 Q - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,; T2 r% O. J3 ]1 V1 w6 j
- * 实现针对特定uid推送数据( v( M' M$ ]% ?- V$ z6 N
- */; _5 N! P* O/ G: R+ `8 |$ @
- $worker->uidConnections[$connection->uid] = $connection;
2 S) u: ?& f8 s. c, S5 J% S* V - return;
- r; x) F( B( {% }6 f' W" q - }! D' F& k E. ^8 j4 T+ c3 j( ?: V
- };' U6 I4 L( K+ I7 X- b1 ?$ {( ~
- # `& p" O6 n; X- U" S& Q j$ D
- // 当有客户端连接断开时; C, C# U# d" _9 S
- $worker->onClose = function($connection): G2 c% q! ]# E$ H/ j8 u
- {
" L, O) x$ T. N6 L5 n - global $worker;
V* T: k, B, y0 @3 Q' Q - if(isset($connection->uid))7 k' f, Y0 M( q4 Q! s" X
- {. N9 y, L6 c/ m2 H: _/ T
- // 连接断开时删除映射/ [, z& F5 E2 u" {2 Q4 t& k2 ]
- unset($worker->uidConnections[$connection->uid]);
+ x6 n3 c+ f! D& r - }
# n1 u$ \1 U2 L9 C( B9 v9 S - };0 }1 b# P4 A% `, j( m5 ? ^+ M
9 A/ h7 X7 r& @/ ^- // 向所有验证的用户推送数据1 r' O$ U) ? B$ j+ @7 f
- function broadcast($message)
( T5 U4 x- j8 i! i8 O - {
5 z# Y2 }5 H4 _- T& K2 u - global $worker;* A7 T. ?& D2 m5 x' n9 l: ~% U
- foreach($worker->uidConnections as $connection)+ v: {! D7 _+ M
- {1 D" F# f/ J0 ] j1 } u1 K
- $connection->send($message);$ p4 p( W: W6 h0 Z, s
- }/ {3 [' ~& ^/ G- I j$ D8 q- L
- }
$ H6 {2 ?9 R3 @* O5 | - 4 E8 ~1 S- K' E2 O* K) v( `% V; z
- // 针对uid推送数据! l) H7 ]4 J, b/ j8 [/ i
- function sendMessageByUid($uid, $message)
# m& |. Y+ I' G5 n% v! e3 I# K5 Z - {
" o7 |# j$ h0 ~7 ?' K - global $worker;
6 h* a4 d( j) [: M - if(isset($worker->uidConnections[$uid]))
2 Q( D- y5 V$ S U9 ? - {
) Z! V# H3 c, w - $connection = $worker->uidConnections[$uid];
. i3 \2 _+ J) y. _6 b - $connection->send($message);) w5 X, H+ C- E2 a; I
- return true;# n, h Z* c @' k; M0 ~7 o: d! {. n
- }
T( {: j! A; ` - return false;3 C8 Y% {! j8 V2 u
- }
v' x' G' d W: t; ]# }
) }) H) E% K# u7 V' k' |- // 运行所有的worker
, A2 e2 V' s: g) ^1 z' s - Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');" n6 K7 i5 [! d6 c0 w- w( L% G5 t
- ws.onopen = function(){
6 T/ Y5 v+ G5 l+ ~& q2 G, ? - var uid = 'uid1';
1 T$ I- j( Y1 `$ k2 |3 F1 { - ws.send(uid);, d: \0 w o" T: |& T8 P3 M
- };
: I) u; S5 |2 U6 D3 t - ws.onmessage = function(e){
. O, a4 v" P. o) e4 k) Y9 c% z - alert(e.data);# ~5 W+ V7 B; ^9 m3 q, f: ?( F& p
- };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
- S8 U' |1 s' S% _; ~7 l% w, O - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);$ t. e* ]4 E. x# n% R& F
- // 推送的数据,包含uid字段,表示是给这个uid推送
$ h# v' k& E2 ~( N+ f: z - $data = array('uid'=>'uid1', 'percent'=>'88%');
# q- ]% _$ W2 \/ d/ \% ~ - // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符# V! R* U* ^& V4 b, G7 S
- fwrite($client, json_encode($data)."\n");
- Q4 n4 O8 S8 a/ o$ h - // 读取推送结果
+ s: |+ x- v/ D - echo fread($client, 8192);
复制代码
4 Q9 Q- q& X- [1 m% y
. d- g z2 c) c3 M' m+ S- I |