- 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;
: u6 d+ w8 t# o b+ x- E& e - require_once __DIR__ . '/Workerman/Autoloader.php';- b$ @/ o: ^4 X
- 1 s6 i( e$ t! l3 H" U
- $worker = new Worker();
1 B( ^9 H* \$ U6 X1 O - // 4个进程
& C8 J; j( ?" s7 ~, s - $worker->count = 4;
7 Y6 |0 | u0 s% t* g' L - // 每个进程启动后在当前进程新增一个Worker监听
' |3 F4 B$ D7 Z6 Z - $worker->onWorkerStart = function($worker)
1 B9 O+ h/ L8 c0 I - {
5 c: `2 U9 S9 c5 \. d, p: w! I - /**
7 |# g2 b$ w8 B; ? - * 4个进程启动的时候都创建2016端口的Worker
0 u# _$ ]5 r) n; h r8 Z5 s - * 当执行到worker->listen()时会报Address already in use错误0 W A# K I" C3 t3 o
- * 如果worker->count=1则不会报错
9 \( s* `+ ?9 b2 E5 p - */$ r* x/ q% w" O2 X w
- $inner_worker = new Worker('http://0.0.0.0:2016');
6 y7 l& p \& U$ ?/ M - $inner_worker->onMessage = 'on_message';8 p) F, F9 ]- v* x( x# W
- // 执行监听。这里会报Address already in use错误
( ~; u+ P+ X/ ~: ^- c - $inner_worker->listen();
7 I; Z' t2 @, t7 ^ - };- t3 a; |) A! Y: M
- @* f. U8 Z V5 z, |/ j- $worker->onMessage = 'on_message';
: h, H- K+ ^3 O/ z: z3 g
( Z5 Q& k# E$ u5 L- function on_message($connection, $data)
1 V# k- q* Z( X' B - {
" y4 ~# @2 o. k; K$ S" G - $connection->send("hello\n");# {3 o0 M3 @6 o9 x. V
- }
1 I9 x, R/ I! G4 \1 q5 g4 W
3 r- _- e. P+ J- // 运行worker
) `! g6 s9 Y5 q - Worker::runAll();
7 u. p3 R! e: |0 a6 U, F - 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
9 n1 a+ y- m# y& x# J - - P. I' ^7 L9 O8 {
- use Workerman\Worker;: B8 I% U8 _& o8 G
- require_once './Workerman/Autoloader.php'; L V7 L t# k; y
5 \& V( P" e- m# V' ] @- $worker = new Worker('text://0.0.0.0:2015');6 s7 ^; N( F6 ^' r3 `: ]
- // 4个进程
, s$ J( a; k3 b - $worker->count = 4;
- ]$ W1 d- R$ V1 ?' l - // 每个进程启动后在当前进程新增一个Worker监听
7 H: P" `& Z2 Y$ H$ C* s, E - $worker->onWorkerStart = function($worker)( L+ Y% ~0 P; T. J4 u8 p
- {
. l9 R' h/ d1 L& `6 A; l - $inner_worker = new Worker('http://0.0.0.0:2016');
3 P! ]2 P* C) M9 @3 O/ Q - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
B) y1 w& M' c+ s- A: G - $inner_worker->reusePort = true;4 o; w, w' Z* \4 {9 t7 V6 q1 V# v4 S
- $inner_worker->onMessage = 'on_message';3 h- M ?8 j6 ?. V' { o
- // 执行监听。正常监听不会报错
! U2 |& g0 Q0 h: o - $inner_worker->listen();
" l- e' F" ~# L3 ? - };
: b2 f1 _( ~( {$ l5 A, g6 N4 p+ k- H" B - * p5 ?: o- h* K0 F# T8 H
- $worker->onMessage = 'on_message';
% g) y: ?4 r5 E5 B
0 y3 F% h# ^& n( u9 D! D8 o8 C- e- function on_message($connection, $data)" ~% ~" T5 f+ m# X7 D" m1 G. J5 N# A
- {
B& \+ {9 g' Z3 q/ C5 N - $connection->send("hello\n");- P8 ]2 B* V1 } C, p
- }+ t; \& {$ b9 j& ~ z1 Z' s( e+ y+ u
- & j& ^! F1 A" r4 f) n, x8 L
- // 运行worker' x s8 A% v( ^% x% l- |+ 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 - <?php8 h" A+ E4 Z8 Q, U7 g
- use Workerman\Worker;9 \4 {$ w) y2 G: Z& [
- require_once './Workerman/Autoloader.php';
, N/ n7 ?2 O/ k7 D& K$ k - // 初始化一个worker容器,监听1234端口/ F" n4 N% s2 `+ o
- $worker = new Worker('websocket://0.0.0.0:1234');
1 w: O2 W# h. ~+ ?* L - " E2 F% B- w5 D8 r
- /*
0 Y: \9 f4 s1 J$ r# n - * 注意这里进程数必须设置为1,否则会报端口占用错误7 s8 r# l& @ F. p; C
- * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)9 L2 U7 R+ n3 M7 @& u( H( W
- */
+ U3 ]1 p+ t2 z6 | - $worker->count = 1;
, A2 e: u% j% z: |- x2 Q# j( k - // worker进程启动后创建一个text Worker以便打开一个内部通讯端口! s* S( K/ ]6 Y5 q' o
- $worker->onWorkerStart = function($worker)1 o* i$ T0 F% P, P: N
- {- C& E9 @: J0 b, E/ d
- // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符3 _- v& A* k* p" ] w3 T* i2 k- @
- $inner_text_worker = new Worker('text://0.0.0.0:5678');
3 A) y0 |2 _' c0 l$ K - $inner_text_worker->onMessage = function($connection, $buffer)
$ t/ q y: V- ~" X& y - {
* N/ ]) M! f" j4 ^9 S& P8 M - // $data数组格式,里面有uid,表示向那个uid的页面推送数据+ n! L2 b) B. T" K6 t8 h9 M
- $data = json_decode($buffer, true);* e0 i. i' a& V4 s; _
- $uid = $data['uid'];
* c: G0 B5 z+ ^ p. Y: E" t; D( Y - // 通过workerman,向uid的页面推送数据
$ k/ ]- `* l/ |# P+ v5 y - $ret = sendMessageByUid($uid, $buffer);
! t; ^) s0 N3 Z8 ~/ A( E: a - // 返回推送结果2 _7 \. ^2 N5 [7 d; c6 q ~9 A; a
- $connection->send($ret ? 'ok' : 'fail');
' T+ S0 a3 V* M5 Z& [# \ [( z - };
6 N9 j7 q: c! p; G/ Q# j0 B. z0 K - // ## 执行监听 ##. d" n. \8 O( i1 @% G$ ~1 ?
- $inner_text_worker->listen();7 B: ^6 r% N5 k* l4 H8 k2 @7 J$ \
- };
& [5 c k8 Q9 o1 G+ G9 y - // 新增加一个属性,用来保存uid到connection的映射
- c L' q% {' b) R, W8 s8 b - $worker->uidConnections = array();( q( V, c; R! Q: X+ P! L3 [
- // 当有客户端发来消息时执行的回调函数
0 v; A W+ _- B0 A, C, p - $worker->onMessage = function($connection, $data)& v/ H) o( k- }
- {4 r2 }+ J# G; S) }& c
- global $worker;
1 l$ c! Z1 `4 ]# ^* y8 { I% P: F - // 判断当前客户端是否已经验证,既是否设置了uid: D3 a" d" `- p! C
- if(!isset($connection->uid))& X! g6 f& f5 M
- {, p" X: X: q. V' y: r! o
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
0 P+ G3 g, X% F! J5 H, z - $connection->uid = $data;
* J2 p; u9 `, s, m9 |7 f - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
F/ b8 b5 w1 G/ O - * 实现针对特定uid推送数据
- R6 u' Q0 a+ {! S$ q6 o - */
* w* Z- k- ]3 C - $worker->uidConnections[$connection->uid] = $connection;
6 B5 Z+ w5 h: t! ~# d! r. w' G& V - return;
3 ?! q! N) O" t. J - }2 ~5 i" Z- b: F5 M9 h
- };
# O& Y9 U+ U) C# K x; |. F, G% j - + {+ z6 }: V. m! K ]
- // 当有客户端连接断开时
3 _( _) r+ J% }& F1 _5 X - $worker->onClose = function($connection)' N ^9 @4 R/ J/ i
- {
% |' g; R& |* R' P/ z - global $worker;
5 e) `2 _2 O- U" T) }% M - if(isset($connection->uid))3 S# v9 f0 X; L1 N+ x' u
- {5 _" F, I- W. J5 t+ v
- // 连接断开时删除映射) h% X' C1 T2 [2 d5 z, |
- unset($worker->uidConnections[$connection->uid]);
! t. D9 U. n5 o X - }2 b( \5 K1 K) ]+ l
- };, k( k: N- H1 L- i! g: [
- ; p# m6 t4 T0 V' S4 ]3 X
- // 向所有验证的用户推送数据
$ U7 Y1 r) {4 G9 T/ | - function broadcast($message). o* g$ K& R9 C) N8 r* |
- {
4 {# R/ e6 Q1 I. e) T: b# { - global $worker;, ^1 `" d( h% Y, f( E, L0 l
- foreach($worker->uidConnections as $connection). F( v0 D% h. [" Q
- {
, h6 W+ c: _+ x+ |) r - $connection->send($message);
8 v7 c( d4 K% ^! O4 U - }# ~ `. |; C6 B. U) h% F
- }& E; v: E+ _3 x$ I. J
! b( X3 e" C& ^2 }# b- // 针对uid推送数据
' V& A; `0 J0 j. u9 Z, { - function sendMessageByUid($uid, $message)
2 V: a! i, r* }& S' L' F - { B5 x2 J+ B" I" f5 f- R1 k; @, E
- global $worker;( S: s4 r: F7 T! S' c
- if(isset($worker->uidConnections[$uid]))
* n0 p- d7 o: ]9 A - {
8 i/ Z$ F% H3 {* v j0 Z/ Y - $connection = $worker->uidConnections[$uid];
4 R) r6 g, ]6 j0 p$ Z8 d1 A - $connection->send($message);2 O$ k& L& s+ I8 X& U
- return true;
- n: r7 D3 U/ A: H6 ?% J - }( r& R+ Q8 M: D! C( \$ R
- return false;
' s8 ~0 _- d! _2 y$ P" M1 p, [ - }
0 {$ U; O, K+ q' m
6 ?! z; I8 p2 a& x- V) A0 Y0 |- // 运行所有的worker: W2 d* m! G! ~4 a7 G
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');# @& N( O `7 k; P8 l
- ws.onopen = function(){2 Q7 |* _4 A0 p
- var uid = 'uid1';4 v: O- f8 b+ s, |6 Q
- ws.send(uid);
: \4 V) ]" O8 h( Q% U - };
, X% L* C/ `8 s J( f: y7 M6 w - ws.onmessage = function(e){; `1 u; Q, m2 `$ E, t) p9 l7 Q
- alert(e.data);
* e: Z3 \; K' `8 |$ S# C% e% @3 `( n - };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
* p) q' d8 c* Z7 V. E - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);4 V7 p5 R3 p% X" Z
- // 推送的数据,包含uid字段,表示是给这个uid推送
, B. ~: T. w+ ?; u - $data = array('uid'=>'uid1', 'percent'=>'88%');- i- o$ d' N; |7 S X
- // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符# {0 i; [7 x$ l" ]4 i/ z
- fwrite($client, json_encode($data)."\n");2 I: A4 ]0 t% [0 ], `7 r% y
- // 读取推送结果
% v; R v% q1 K9 T) j e4 ^ W - echo fread($client, 8192);
复制代码
4 Q+ Q1 \" [ d+ X# y2 z7 F! G/ ?9 C/ Q; l& f
|