- 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;0 N. g1 E+ Z. s
- require_once __DIR__ . '/Workerman/Autoloader.php';
% V) V+ @& _: r& C3 u3 j
5 O( r# u5 [) v: j1 S& w& ^- $worker = new Worker();
: u, a# o- U3 e" m. ]5 U) \ - // 4个进程
2 J& \# C; I) }( E$ C - $worker->count = 4;/ t* I+ g7 M; v/ L
- // 每个进程启动后在当前进程新增一个Worker监听
3 |: p' ^7 o$ M - $worker->onWorkerStart = function($worker)
% G' Y+ Z7 v& e/ e - {
{* E0 r0 z/ v+ ~ - /**% J/ \9 K: `2 i
- * 4个进程启动的时候都创建2016端口的Worker
3 {! r, f7 K; @/ A. C - * 当执行到worker->listen()时会报Address already in use错误- p8 `" p. x7 U- f! ~" E8 _& K [+ V
- * 如果worker->count=1则不会报错0 l/ Z- f m. _$ A' X
- */$ X4 }) M9 l, z, f& B
- $inner_worker = new Worker('http://0.0.0.0:2016');
. e) r! @8 o% h' N7 D, G4 y - $inner_worker->onMessage = 'on_message';
4 q1 H6 I5 ]/ f1 |- R0 ?% j - // 执行监听。这里会报Address already in use错误; A9 K/ a5 r# ~- M1 a- w
- $inner_worker->listen();) C4 I8 E# d* |8 _& g
- };
2 Q# W8 l" p! e H% n& ~ - / U' k- r4 |) [: {# R5 w
- $worker->onMessage = 'on_message'; x2 w& Y& ?3 y% r! ]0 ?
& a: X1 F! q3 j; ~- function on_message($connection, $data)3 v3 t" g/ ~. s6 @2 T+ o0 e9 K2 Y. M
- {
0 \% N' j1 f' z - $connection->send("hello\n");
& G4 v. r3 ^. I - }) T4 {% S/ u5 e
- % X1 N0 N1 b/ ^9 n- O
- // 运行worker# J- ^/ l# F2 k9 X
- Worker::runAll();
% f/ w# V7 ]7 g' ~' O - 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:' L7 W; m+ Q5 J- M
& @6 _4 D: T$ c% F ]+ O1 ?3 P- use Workerman\Worker;
. S4 p. D) Q# H7 [1 l( P - require_once './Workerman/Autoloader.php';; k% n+ Q3 Q6 ^6 [
6 F( ]8 J7 N R- Q) ^3 E0 ~& o$ E7 `- $worker = new Worker('text://0.0.0.0:2015');
# \( P4 K1 o# L$ Y - // 4个进程
9 i6 {6 f, R/ b; q - $worker->count = 4;, q; W0 z Y) s% G) x
- // 每个进程启动后在当前进程新增一个Worker监听
H9 T* b' e3 J - $worker->onWorkerStart = function($worker)# ^& x8 A' M8 e& c* v5 e. K
- {5 t7 Y: T4 T% U7 g; m
- $inner_worker = new Worker('http://0.0.0.0:2016');
7 x+ o8 y4 v& N) C! A3 j7 p - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
7 f2 m0 P0 S2 I2 j s' k; G- { - $inner_worker->reusePort = true;
) y$ W Q5 u" ]$ O+ @/ H - $inner_worker->onMessage = 'on_message';4 Z: k* m! c! K' Z& ~# a
- // 执行监听。正常监听不会报错
/ q0 l2 T. X* d3 x# | - $inner_worker->listen();6 F' k* I4 j2 L& o
- };
7 K7 d: r9 b- k$ s1 Q - ' Q: O8 V% ]/ l- i# n% {) V
- $worker->onMessage = 'on_message';
" D; D4 P9 Z. i8 D, D; M6 {6 a) l' ~1 X - ( X. A# i' Q0 F" {; C, g; K
- function on_message($connection, $data)' l7 w) X- u" m8 h' h# O
- {
! u; q9 S; \! b - $connection->send("hello\n");
& Z# i7 b9 [8 h8 E4 C( B - }
6 e3 I G) I& [' ~: F) \ - 6 [2 f) e# q* b4 r7 s& s' R A# c
- // 运行worker
- r( U/ b$ i( L M1 c: E - 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: i, n/ s9 n' \& i) e6 O
- use Workerman\Worker;
7 _) ]6 A% I0 y8 W& C! z - require_once './Workerman/Autoloader.php';
e' H1 G2 s. {) `! w: r2 Y - // 初始化一个worker容器,监听1234端口
; d+ z4 P; M5 V/ `/ Q - $worker = new Worker('websocket://0.0.0.0:1234');+ o$ v: c- D3 I) ^3 y. G
4 `1 L9 Z! T$ K4 Y: d- /*
) ^. [6 I7 v- T - * 注意这里进程数必须设置为1,否则会报端口占用错误
; [9 K0 U J7 o2 p3 ~ - * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)- b% E E" I( P) g, P
- */& y6 E% I6 T* S% G) J
- $worker->count = 1;
& G9 \5 j! ^. N" R0 t g4 ?" H3 N - // worker进程启动后创建一个text Worker以便打开一个内部通讯端口0 L @( K& k, a% ^* W( Q4 s/ J
- $worker->onWorkerStart = function($worker)
; g8 u u, i5 E& \# D - {
' P' J. x# I% n! x/ w$ m - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
- c5 R) u! k1 l& Y( [ - $inner_text_worker = new Worker('text://0.0.0.0:5678');9 _1 i5 c) h, a7 p2 }6 }
- $inner_text_worker->onMessage = function($connection, $buffer)
1 s% l1 ^+ g8 f! S5 B+ P* \! Z - {- ]! v- w) d/ @: u6 Y* g8 S
- // $data数组格式,里面有uid,表示向那个uid的页面推送数据
. s8 ?) h7 \$ m, c) G" i3 |! j - $data = json_decode($buffer, true);- X& I' R& O& y! v& E2 U
- $uid = $data['uid'];+ I! l% K: Y3 `/ {
- // 通过workerman,向uid的页面推送数据
% E* W0 F6 W3 Z- d9 v$ y: P - $ret = sendMessageByUid($uid, $buffer);! j+ U" M* V) }9 q5 K4 `: {
- // 返回推送结果0 r, { t- x1 |* p' n$ ?" h N
- $connection->send($ret ? 'ok' : 'fail');
# N) c- p b# a4 [9 z6 w% G$ ^ - };
9 B$ z: H& c) I% I - // ## 执行监听 ##
, ]; [: E, f, T0 d - $inner_text_worker->listen();
% ]' S$ P, x) n) }0 M - };6 g9 [) D+ G8 G- U# u6 _
- // 新增加一个属性,用来保存uid到connection的映射
2 l3 ?" S' C6 [- S3 m( N. L* g - $worker->uidConnections = array();
! J: i5 d: w/ \4 O w- i - // 当有客户端发来消息时执行的回调函数( U0 k# k. r# L
- $worker->onMessage = function($connection, $data)
& [. n3 [% p& |% g5 y) v5 f* ~: G - {
7 ]1 Z0 P) v5 l6 `3 D" t - global $worker;
. e2 n! }% Q) L5 U - // 判断当前客户端是否已经验证,既是否设置了uid
, k* O1 p+ V3 b" ?% V, w! Q9 i2 b - if(!isset($connection->uid))' @( k5 B) _' _$ {6 @7 x% U
- {
* P- @9 _4 U0 Z3 v+ v+ G - // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
2 M0 W8 h7 M" j - $connection->uid = $data;* D8 N9 ~* a: Y) ~) I& i
- /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
9 b G8 Z) U- L1 } } - * 实现针对特定uid推送数据
" N( C6 n2 d( ]7 l h& x2 ^ - */
1 S' Z5 f# T8 y! h$ }6 I: h - $worker->uidConnections[$connection->uid] = $connection;
$ {& t {. r/ m8 {, @' d - return;
! i8 W# P" Y' g5 n; M+ p2 f5 y5 m - }: J) A- y5 y4 \0 ?5 T) H2 C
- };
; ]2 J2 b q w6 b# }0 g1 t1 \
8 Z2 E/ \% b+ R- // 当有客户端连接断开时
4 `+ [3 i+ S, E - $worker->onClose = function($connection)7 ?+ H$ n. V& k1 ? `1 O: x9 y2 @- m8 P, h
- {7 E4 o5 X0 W- d/ d1 ^
- global $worker;
* R- Q9 W, `" i0 B - if(isset($connection->uid))) o# F* x5 a$ U0 e6 G
- {/ P2 z) I& M# b" ?3 A9 S6 z
- // 连接断开时删除映射2 z2 e, \/ P9 I4 }. }
- unset($worker->uidConnections[$connection->uid]);
, U E' f* G1 C% m( r" q* P - }
( T8 b* R8 H/ `1 ?& V( F3 u/ z! r - };
G+ W. i5 J0 V" y1 g- d, u* h - 4 D7 j9 D- g* i# O1 r
- // 向所有验证的用户推送数据+ G V* m5 v3 Z& g0 b
- function broadcast($message)
3 u# p1 i2 }+ C" C) @ - {2 b0 X# S! k, k
- global $worker;
, f, U# W1 U% Z - foreach($worker->uidConnections as $connection)4 {, c1 t9 y& s( r8 d9 ]! c
- {) z7 ]/ X c! |0 B! o- P1 o
- $connection->send($message);# t- C9 N: S5 @! q. @" U+ T/ }
- } v7 v) i( j R! x2 k, N" l
- }6 v; ]8 X% d+ `/ i0 ]
- ' i' Q6 d4 p5 H" A
- // 针对uid推送数据' ~# P- m- G" k
- function sendMessageByUid($uid, $message)3 ^7 R! N# h+ h: h; G- F1 @5 F6 u
- {
* `0 N. w$ j# Q - global $worker;
; n% p1 r- G A" D8 }; F" k; d" } - if(isset($worker->uidConnections[$uid]))! C8 P8 I: g- z# Z
- {; i7 N) C8 w d6 B$ u+ Y
- $connection = $worker->uidConnections[$uid];
$ Q1 n7 N- k: u$ s0 D7 M- H* _ - $connection->send($message);
( r, g9 ^! | {. `- {, t4 _4 @ - return true;
F; P- {& g/ [3 C6 m4 x - }
0 [- @: U# Z. X3 X5 ~$ }4 ?4 H - return false;/ U2 i! R+ \" a6 ]8 V; F
- }& v6 x I/ p- F/ `; k
+ L% m! p# |3 S- v6 n9 X- // 运行所有的worker
' l1 r: E* N- e3 p+ k0 ]2 w0 a - Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');
. _2 a ^2 [" _ - ws.onopen = function(){
' E3 t: R3 [- O$ }( [ - var uid = 'uid1';
/ q' y. ?& {4 b - ws.send(uid);
9 w+ t9 S' s$ W4 n* m6 K - };
: d( G+ g( z5 P S3 y9 X - ws.onmessage = function(e){2 M" W5 D; Y5 S/ ^ }
- alert(e.data);
, A: J/ p8 Z5 L6 ~ J - };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口, \$ q/ ^4 A% }- E2 f B9 z6 v; h
- $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);
% @" L. G9 n9 q3 G/ D - // 推送的数据,包含uid字段,表示是给这个uid推送- h' O' |, {; ? A
- $data = array('uid'=>'uid1', 'percent'=>'88%');
/ D6 B# }# v4 }, T1 I* g - // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符" k$ m1 Q# b) y) g* r
- fwrite($client, json_encode($data)."\n");
. F& ?6 d& s3 g9 N2 A - // 读取推送结果
0 Z% K. f4 r% v! V, | G* g* l& {9 | - echo fread($client, 8192);
复制代码 & m5 R: P7 s3 V: N6 i- {7 g
( F7 ~0 j0 N3 e- C r w- _! B |