- 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;
) x: a: X/ _8 w3 J - require_once __DIR__ . '/Workerman/Autoloader.php';
! C/ b) S! [# j( `% b- V - 1 ~+ e3 c2 [$ `; }
- $worker = new Worker();# F9 z( I4 z. _& R7 G
- // 4个进程$ L2 O" d* e. [
- $worker->count = 4;& ?- g* P& U5 W3 _2 m
- // 每个进程启动后在当前进程新增一个Worker监听
) E$ d1 z( k- o - $worker->onWorkerStart = function($worker)
& h F1 w/ G2 `' O- @; o0 w - {* ]1 o) O. n9 W# A' @
- /**
! y- O8 h- O9 W" w( l6 f - * 4个进程启动的时候都创建2016端口的Worker
9 H# a9 N6 q- D* V7 g k8 w3 B - * 当执行到worker->listen()时会报Address already in use错误# W# Z8 \8 K( K7 E- h6 s# h8 X7 z
- * 如果worker->count=1则不会报错( v- H0 a* J" n
- *// J8 o2 B* o! T* ?' U! [
- $inner_worker = new Worker('http://0.0.0.0:2016');& C. I+ m4 Y" l, |
- $inner_worker->onMessage = 'on_message';
( ~2 x- q P C) r5 e - // 执行监听。这里会报Address already in use错误9 A1 L" j/ G+ v) S; t
- $inner_worker->listen();
& d# P# }& B0 }4 l% g - };
" N8 @ q2 E2 I* e3 [9 T - 4 i% |. X; V0 E. F
- $worker->onMessage = 'on_message';
$ R* g# a; y: A2 a2 I. p - 6 H Q$ L G' J. x& k
- function on_message($connection, $data)$ v' c8 Z/ C4 c m4 f
- {+ r- S8 M8 J% D; l- I) k
- $connection->send("hello\n");% y0 p, e, F( i7 _: N2 d
- }
' c- j/ l3 A: H: A
' V7 D: K/ r! w. x- // 运行worker
1 I) k2 \& J* V: F) e3 G - Worker::runAll();8 N, z2 _5 @( s3 \% F8 C- S# U a
- 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
( O5 B, {) H: p - . t4 [! S" X# |+ ]1 g8 o, R$ E
- use Workerman\Worker;& ]( h: i ]4 A0 V. B6 T
- require_once './Workerman/Autoloader.php';( o, q1 Q; x. z' U8 b' Z# y
- Z) q P. }( K! L8 D5 N- $worker = new Worker('text://0.0.0.0:2015');
, F& }/ R- R9 b w& O* U" e - // 4个进程
7 R$ W% ?6 ^! c1 R) K1 Z- d - $worker->count = 4;
( B3 S7 U2 n, f) N& ] - // 每个进程启动后在当前进程新增一个Worker监听: Q' I6 d$ j8 \6 D3 n9 R8 k; F) E
- $worker->onWorkerStart = function($worker)
2 V% F: a% r4 s1 F5 R - {$ g/ g* e; o! d( ^
- $inner_worker = new Worker('http://0.0.0.0:2016');
1 I& \8 G) y( W! C; Z% {; B* h - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)$ ~! {0 L5 s5 M3 l
- $inner_worker->reusePort = true;: X3 R! z1 U* ^# s6 Y
- $inner_worker->onMessage = 'on_message';' F& {7 O* f* u/ X
- // 执行监听。正常监听不会报错
. E1 R$ ?, x+ u - $inner_worker->listen();
5 m; u. Q) F. J/ g: ~ - };
: |, {0 K$ J$ o6 E# M
% q: K+ A6 o0 D1 t$ z- $worker->onMessage = 'on_message'; y1 T5 ~, t$ T3 x- t
0 C3 T9 Y5 ^# c% N2 z8 q" ? N& C% L- function on_message($connection, $data)+ ~% G- d( S: T9 d: u/ z8 Y' Q9 J
- {" Q8 P: q" `9 A' u" u& R
- $connection->send("hello\n");
8 j) w8 ~0 D- R) j) G: X. {+ P - }
5 _ x/ f& c$ H) R - + B6 C, {8 t# a, X, t
- // 运行worker* N; H2 X6 w2 @. A
- 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 - <?php9 S0 w) P+ j6 W+ k$ n/ p0 p2 [) T
- use Workerman\Worker;
/ J8 v! i: V! [. u - require_once './Workerman/Autoloader.php';
4 n% I, o% u5 r* r6 h& E: n8 m - // 初始化一个worker容器,监听1234端口
/ g4 z( \! @$ I - $worker = new Worker('websocket://0.0.0.0:1234');
* f5 P9 e2 j k
8 U; s' n+ e9 ]- /** A1 i/ g6 k2 ]5 D& ^, s h- c" Y8 V- T
- * 注意这里进程数必须设置为1,否则会报端口占用错误! F& L/ i' x; j# _4 D+ ?$ j
- * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)
, c y! X0 l/ S2 i% o% z3 M& \% n - */
! L6 A' c2 x; v( y9 V9 T - $worker->count = 1;
$ t, z/ ^$ o7 e( T3 I8 Y - // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
6 b! c5 K5 e8 @1 P6 T7 u1 m - $worker->onWorkerStart = function($worker)
. a5 I% `7 p: f0 U7 D+ w3 Q - {& ]7 c8 l/ [3 D- \
- // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符$ a& ]$ c1 |* b/ Q
- $inner_text_worker = new Worker('text://0.0.0.0:5678');5 t1 c+ t/ I9 b E) T% H; N5 J" V
- $inner_text_worker->onMessage = function($connection, $buffer)/ R" ?/ E5 }4 o! T$ W _/ s
- {
; @/ v$ E7 W8 R# b' G - // $data数组格式,里面有uid,表示向那个uid的页面推送数据! U+ h7 f) m! G% V+ n1 w
- $data = json_decode($buffer, true);
5 I8 t5 ~8 M }" }: q2 V9 X6 A/ n - $uid = $data['uid'];
8 P: X1 S/ g4 V. Z6 g/ s2 U" G - // 通过workerman,向uid的页面推送数据
- M: t( G/ H F2 g) p - $ret = sendMessageByUid($uid, $buffer);3 `4 h1 D: ?# B/ C$ r( s
- // 返回推送结果
& A, Q7 l7 ]( N' S - $connection->send($ret ? 'ok' : 'fail');0 l3 U1 c7 N, Z
- };
+ k4 X) F" c2 S4 K- e - // ## 执行监听 ##
`& X5 T- T6 K* I1 |) d - $inner_text_worker->listen();
5 l$ P5 u. _- o+ v3 D: F! c- p - };5 v. j8 K9 k1 W8 V; l
- // 新增加一个属性,用来保存uid到connection的映射" } |' s- K! z# D+ G
- $worker->uidConnections = array();
9 K; P$ s2 y- }: F) q - // 当有客户端发来消息时执行的回调函数
, [( B5 f8 K$ @2 H I# W - $worker->onMessage = function($connection, $data)) @; y- J; o) n# X4 Z) {; g
- {' _ T' y% T) t1 v, q0 H9 L% g' D
- global $worker; b- h9 B* F( K- I' ^
- // 判断当前客户端是否已经验证,既是否设置了uid
( S9 B" o! D& r% B/ X& H' J- ~/ ? - if(!isset($connection->uid))
5 I; N* U% f9 G- ^ f$ H0 z( P) D - {! R7 J. |; ^( }
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
) A# k/ G, ^7 w/ n - $connection->uid = $data;
& Z: g0 K& t# w) A& t - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
- b* z3 [$ G" O( \8 d5 _ - * 实现针对特定uid推送数据
3 X4 K" h+ J; O" \1 m. D6 B q% j - */5 B' }+ d) W9 R* y% M
- $worker->uidConnections[$connection->uid] = $connection;) ~( P# U; `: Q1 O T$ a1 x8 f
- return;
, C R4 @9 S7 @$ W9 ? - }& T/ B2 t9 v6 E# `$ S7 O0 ?. {
- };! f6 y6 U2 |5 X4 d
8 n V. w& [% J j- j' r7 P3 H+ F" B- // 当有客户端连接断开时
8 F1 @1 m0 X& r6 L - $worker->onClose = function($connection)" }7 C- L |) k" X* H
- { B" P% L# k# j
- global $worker;+ h1 C; @! h+ q) d6 p0 O1 X
- if(isset($connection->uid))0 O; \6 J6 F. r, U7 k. s" M5 E
- {
/ ]; t% @5 X% y0 G' f/ O5 I* G7 V - // 连接断开时删除映射
( B! Y1 B$ }/ I3 Z6 I; q2 x - unset($worker->uidConnections[$connection->uid]);
5 |1 P: G# B j! G0 L; q9 G - }
2 L8 y( c2 c" D - };
, z' F1 }( t! I" t5 H
" n' p: {9 V+ d1 h$ _( ?0 `, r- // 向所有验证的用户推送数据
4 a' I R& v0 s. B4 I$ } - function broadcast($message)% i7 D. g# A* V T. h* B7 m
- {. Q m& M( C9 ?- b% i
- global $worker;
" c) Y: L2 I, }* H) A, r j - foreach($worker->uidConnections as $connection)
% \( ]2 T6 D v& A$ @8 O# K - {
9 K; r. y$ x& |' n! | - $connection->send($message);
4 Z f" U. j, K; g - } z1 o8 j! m2 U, F6 q
- }- t3 Q4 k/ U6 s/ V1 q/ p
' H5 @% j, K+ t" s3 }- // 针对uid推送数据6 E2 T3 ?; h6 L- F# C/ K5 q' k$ r
- function sendMessageByUid($uid, $message)( A' P6 r! q N: @
- {
. I3 l& \2 H- m0 s! l- f - global $worker;
* c2 w* Q; d0 q - if(isset($worker->uidConnections[$uid]))
* S3 Z5 y7 j% V - {
N5 N) `% |% u9 U+ D+ l3 Y - $connection = $worker->uidConnections[$uid];
1 { R' ~/ {6 d& X- P0 U - $connection->send($message);- T3 i6 |0 D5 C1 ^4 `
- return true;" h5 h9 {& P! c, ^
- }( H2 P" _' Z7 j
- return false;
1 l0 m' b* \( r3 |0 z - }
6 s+ h3 b1 d3 I# ~2 P! B; O5 A - * ?/ I$ n3 g' u9 O' n
- // 运行所有的worker5 d: k2 @ i9 @
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');& l9 k4 ~ n! E! e# ~
- ws.onopen = function(){
: O h3 W1 v6 r) ~ - var uid = 'uid1';5 R8 s8 [# x9 x n
- ws.send(uid);
. |7 y2 }$ p |* m7 j \7 Z2 v - };7 {" q8 N' z( c) Y" X
- ws.onmessage = function(e){
* n% Y+ X3 D+ O1 P, k, E* s - alert(e.data);3 Y4 q: w+ W4 A% e6 F
- };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口 U: g6 y- s* x! t+ o. {% a
- $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
* C+ ~2 A& I+ \- X$ W3 ] - // 推送的数据,包含uid字段,表示是给这个uid推送6 ]/ j+ _) w ]4 E9 X- K; }9 O9 E, a
- $data = array('uid'=>'uid1', 'percent'=>'88%');5 N3 N- J2 Y( D8 `8 }9 Z* R4 k6 s' |3 b7 D
- // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
. o$ p% I4 B6 Y) m6 h - fwrite($client, json_encode($data)."\n");. [, w$ p2 T9 D. J; y6 W" \
- // 读取推送结果4 |* e" v, ^" j, w1 }2 a9 `' Y
- echo fread($client, 8192);
复制代码
: W, q7 C% e- u
! G* w" ]' o- q |