- 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;
( `* N$ w/ t2 J4 K - require_once __DIR__ . '/Workerman/Autoloader.php';
' v( U+ \, ]) \* u' B2 ~ - {6 t: [8 e7 }$ s+ G
- $worker = new Worker();
$ C' M) i( H* v# f$ H6 |8 E - // 4个进程
' K$ \ m+ ]0 C4 i5 Q - $worker->count = 4;
3 M) Z$ @: C0 b - // 每个进程启动后在当前进程新增一个Worker监听
2 o! h* s+ q0 D! E+ X - $worker->onWorkerStart = function($worker)+ n) q" Y7 k% F# ]
- {
& n6 n7 Z1 r* ^( o$ R& ] a9 y6 b - /**
. n" s6 E" n0 R9 h/ b! W6 s6 x% m - * 4个进程启动的时候都创建2016端口的Worker3 c# d- H' |4 Y3 t! A3 a% K
- * 当执行到worker->listen()时会报Address already in use错误
q7 i% v/ H3 G/ I K/ H. o - * 如果worker->count=1则不会报错
5 a* Y3 d Z& m' e. m& H: a& A - */
/ @2 W0 W9 O4 k, u! @8 X - $inner_worker = new Worker('http://0.0.0.0:2016');( h% M. S% V1 Q+ |
- $inner_worker->onMessage = 'on_message';
' h9 _& f5 d0 j0 W6 h1 o - // 执行监听。这里会报Address already in use错误
S7 W" e0 \& N( X - $inner_worker->listen();' b; H! ^% X3 X" z7 \
- };4 U# A( }. |- b! A# k
- 0 l' w" }( C& L0 {9 c; J. f0 i! L7 z
- $worker->onMessage = 'on_message';+ c* c4 ~: B, C0 q: i4 d# A" b. L
5 L( k& Z. i. n$ Y. }- function on_message($connection, $data)
, |* e6 Q. O {' U9 c - {9 U# J6 F7 ~* w
- $connection->send("hello\n");
" A2 K! i) a( z0 c8 T - }
; P! e( S* H8 N; Z& Q* Q( ?# T# w7 E
; C% [4 I2 T: b' H+ h1 N0 f8 _9 Z- // 运行worker/ Q- k5 _3 g; Q: M5 w: D
- Worker::runAll();
) F+ O& B1 y$ Z5 Z0 S1 f - 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:0 H l' j% Z# n. f9 ~$ e
- - n j9 q8 i7 P4 t
- use Workerman\Worker;
( F% k' `* @/ N - require_once './Workerman/Autoloader.php';
/ p1 u7 W' i/ v8 c! p - # J6 W* T6 X% F
- $worker = new Worker('text://0.0.0.0:2015');. [; ]# [8 }- `) F+ t8 R
- // 4个进程
2 h7 n N& E6 V$ z& D - $worker->count = 4;
6 w% p: A+ u7 x2 E - // 每个进程启动后在当前进程新增一个Worker监听, n& G/ A% z4 W/ k, J; U
- $worker->onWorkerStart = function($worker)
* f# ?- u$ x# u& I8 u I - {
! I" T4 u, j9 L& h - $inner_worker = new Worker('http://0.0.0.0:2016');
# S2 g. ^5 t7 T9 p- t. q; y - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
: J. ]) e- Z/ y O. }: } - $inner_worker->reusePort = true;, g: M( v. x4 u+ Z, Y' y/ s4 I
- $inner_worker->onMessage = 'on_message';
% J* A0 G. u( m) w* p- ^ - // 执行监听。正常监听不会报错) |. ?0 K3 t. f. [7 Q
- $inner_worker->listen();
; N/ v, k6 H% M - };
% u6 T, M! i. ~: d A
7 g1 F7 v" m. z+ ?: t; n* A2 m- $worker->onMessage = 'on_message';% `; m/ S: S! h% ~) D+ V
% m. U9 y' f) c- function on_message($connection, $data)8 N' `3 {( {: s6 W' Q2 |% D* h
- {4 X9 e2 x/ c( ?: E
- $connection->send("hello\n");
% f0 x4 B r! n5 R - } \, Q" P& ^6 q, A4 o5 ~; d
- 2 s6 ?% v% _2 `6 @" ]' k; Q
- // 运行worker
9 {5 E5 y" U* z( I$ W0 h - 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
8 |& {( j3 ]5 Y# |* J - use Workerman\Worker;
$ u( ]0 @4 [! M4 s0 p - require_once './Workerman/Autoloader.php';
" J: S4 O8 M( q4 {' H, h/ i" t( b1 e - // 初始化一个worker容器,监听1234端口
- y9 A3 p0 q' r5 E - $worker = new Worker('websocket://0.0.0.0:1234');
7 v& D, t; c9 p" f% h& ~% ]; Y - : x/ A1 Z/ Q5 w! x
- /*4 v v- ?$ I2 ?! H: B b
- * 注意这里进程数必须设置为1,否则会报端口占用错误
0 f: l3 \5 p6 l$ G& X5 { - * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true) ~% g+ l8 }1 ]0 f; w8 _: d8 r
- */* H6 M, u( {* u3 ?; Q
- $worker->count = 1;6 j( F, u) Y l5 G* w$ } c
- // worker进程启动后创建一个text Worker以便打开一个内部通讯端口5 T# v8 w1 F" f2 C9 C9 L
- $worker->onWorkerStart = function($worker)
: a) F) L0 x# B+ K - {
+ g2 ~2 X/ _( E+ |. v - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符1 x2 S, l8 X: g6 W: W. Z" C$ G
- $inner_text_worker = new Worker('text://0.0.0.0:5678');" s* E2 |& G! k/ Z8 e0 ^" f9 @
- $inner_text_worker->onMessage = function($connection, $buffer)
5 V2 m1 e" j7 f! R. C: V - {
3 Z. T6 P. u! k W, m& o - // $data数组格式,里面有uid,表示向那个uid的页面推送数据% A9 D0 O1 a; e
- $data = json_decode($buffer, true);
5 R3 ]6 r# \ N7 b; u - $uid = $data['uid'];
9 [% F4 f9 Y q' v: c - // 通过workerman,向uid的页面推送数据
, n& I! Z ~" t; D; E - $ret = sendMessageByUid($uid, $buffer);% r: A# F" a0 D& p! `
- // 返回推送结果
' ^# j( e9 p( ~# G4 G - $connection->send($ret ? 'ok' : 'fail');
* l7 u: {6 Q9 S0 P* b' C - };' b( G$ y% B. @% h, e
- // ## 执行监听 ##, Z! \$ p, M( b5 ?! s- H
- $inner_text_worker->listen(); e3 b/ C$ L U; g$ N s' n
- };! V9 u- c/ d) Y W
- // 新增加一个属性,用来保存uid到connection的映射
/ g8 B) k+ \2 R! e - $worker->uidConnections = array(); }9 c: f5 j1 L! x! C: c
- // 当有客户端发来消息时执行的回调函数
' q: W4 Z6 i# C+ g7 v' i5 N5 z# [ - $worker->onMessage = function($connection, $data)% l% C% ^4 p5 A$ G
- {, D( e2 f" r- J- g
- global $worker;/ Q7 R% q# {0 D: ?$ w
- // 判断当前客户端是否已经验证,既是否设置了uid
8 Y H( |5 L3 z0 x- {. r7 O - if(!isset($connection->uid))
, {- C' \" P C& u9 n0 W% I - {4 g, _5 _" T7 z% R% S( D3 c
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)# [6 B9 T& U. h) c
- $connection->uid = $data;
2 ~7 Z9 ^8 u" O* [: [0 Z6 Z - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
; u% O/ w0 k7 r - * 实现针对特定uid推送数据2 y9 Z. e. W$ E( K
- */
6 N3 O# G' I4 S5 n9 W9 h' Y0 U - $worker->uidConnections[$connection->uid] = $connection;
+ v0 v6 f- M4 H6 w0 Y - return;
5 f% e, E7 I, j0 R- W$ ] - }. J# y9 \# B% w% }" h8 Q
- };/ o9 l' r* h. L1 m0 F5 P( O/ O
: k; v8 y1 K' _! j- // 当有客户端连接断开时) ^( m) K1 a6 f! }/ D
- $worker->onClose = function($connection)0 Z8 t" o+ w) [! I* j- O
- {
8 K' P" }0 s4 b% K6 S, @ - global $worker;7 O; ?9 t3 \! a# g/ G
- if(isset($connection->uid))7 z6 w" [ B" X9 {3 [
- {
) l2 o2 n& f9 @3 A, J/ \ - // 连接断开时删除映射% h4 r0 u0 l+ q$ O2 m; C" _
- unset($worker->uidConnections[$connection->uid]);6 h1 |9 A6 U9 r% i
- }
4 R `" B; g6 c" v2 `! G - };6 X C) Z& n0 |) C2 Q; v. Z
- 3 V% y- C; K: V5 X4 w% d7 Y
- // 向所有验证的用户推送数据
$ q6 X0 p* C3 a: x# r( m5 ?& r - function broadcast($message)7 r7 d y' |9 I7 d+ i, c
- {& a! V x3 i1 l7 o; {; u
- global $worker;
v# S; i# E9 U1 R$ ` - foreach($worker->uidConnections as $connection)
0 Z% O' k* L! d1 L' W - {$ p+ ^! I2 o1 S$ S2 \
- $connection->send($message);" f% F/ t; H" M
- }( Q% q }0 `+ c2 y/ }; `
- }' @, V( ]$ v# G* D
% l% U1 J$ o3 g) b/ h- // 针对uid推送数据
3 e* o0 B9 O2 G; {, Y) S# { - function sendMessageByUid($uid, $message)+ A( L& @0 M0 h
- {
; w" G9 [+ T# E3 n2 N' ^7 { - global $worker;& L% b9 ^: V! |) ]2 }% h
- if(isset($worker->uidConnections[$uid]))( ^3 Q8 h, v* }, C
- {
3 l# E& f% z' c. S* g! i3 B - $connection = $worker->uidConnections[$uid];
" v" Q/ |9 {) d; G- N# [ - $connection->send($message);6 z( v F, U& O' P
- return true;
5 b9 g% K @: f+ i1 j; d+ A - }+ N' u' {" E0 ~5 J2 k6 C* o
- return false;
5 r1 R7 E+ w. j1 b" \ - }1 B1 P o0 i( P. ?4 k% }
! @; Y0 h, n- _% t; T( D- // 运行所有的worker
1 _2 `5 P) h" A+ Q! H" | - Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');- M5 B L7 u, e1 b+ t! J5 t
- ws.onopen = function(){
+ a4 N( S' {; U1 ] - var uid = 'uid1';! _1 [# }# W. Y# y. [+ R/ a
- ws.send(uid);
R/ i- @. Y( U5 C4 ? - };
8 w( l& t$ R4 y p0 v8 c' A - ws.onmessage = function(e){
4 R+ ~! a% f* u6 }. f/ K o - alert(e.data);
6 _4 X# r5 T, X% E' } - };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
+ N( [( j8 K8 q - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);* w6 l4 h6 ]9 ]9 Z& V$ j5 i# i" w
- // 推送的数据,包含uid字段,表示是给这个uid推送) N& M+ y0 K5 T" [9 L) Y/ }
- $data = array('uid'=>'uid1', 'percent'=>'88%');
) Q9 P! m$ c; d- b - // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
' a" {0 h. n. p2 g& f: c - fwrite($client, json_encode($data)."\n");
$ f; Z+ [ H+ a+ p3 a - // 读取推送结果
- j, z2 A9 o5 A& J, p+ ~4 S @# t - echo fread($client, 8192);
复制代码
. a0 a2 ~9 X. N2 ?* P6 S# B" }. @
|