- 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;: U7 z& u* ^/ R% o. V
- require_once __DIR__ . '/Workerman/Autoloader.php';
/ \2 M$ i8 j- K; I - : U7 N1 @; C9 x' ]1 [
- $worker = new Worker();" e& Z: M) T8 S) T
- // 4个进程
( Y% V9 i, u" B7 i - $worker->count = 4;
4 z4 `+ q# m- V0 b - // 每个进程启动后在当前进程新增一个Worker监听; K- V' a6 \5 Z- l1 u8 c9 _0 Z. z
- $worker->onWorkerStart = function($worker)' @# k3 `# g8 }- z' A+ Q+ O
- { o" W2 T- _6 I0 R, a Q3 P0 R/ l
- /**
2 G, M, C" m1 s3 c2 L) g - * 4个进程启动的时候都创建2016端口的Worker
8 m! V1 Q' h9 S: F4 \; T/ I$ X - * 当执行到worker->listen()时会报Address already in use错误7 B; o7 F% \* Z8 x
- * 如果worker->count=1则不会报错
: ~- D! t% w$ `8 x4 v0 @ - */
3 ?/ k3 B) l4 J% ~9 R1 k - $inner_worker = new Worker('http://0.0.0.0:2016');
* `, C4 I, U( X4 z+ x6 p - $inner_worker->onMessage = 'on_message';
0 U/ \- X. D) [7 X0 a9 ^1 M& e - // 执行监听。这里会报Address already in use错误
' @ t' V% @/ K8 Y - $inner_worker->listen();5 G3 h8 O2 d, u" \
- };
, `, W, A3 p: o# C M$ h
9 e: a, a. N9 t4 J- $worker->onMessage = 'on_message';( Z) l3 E5 w$ `: i) {
- ' R& S3 D; t! j C
- function on_message($connection, $data)
. t. I, n/ ^1 j - {
* R! z4 j0 N' ~$ N; z2 m - $connection->send("hello\n");
# l0 F. i! E$ k) W8 R& ~! E" l9 y - }
8 Z2 r: _9 S& I( f% k
: ^) Q3 K% K1 K$ A4 j- // 运行worker
' }* {' s1 j7 ^% v1 i! E" } - Worker::runAll();. V( G! g1 i. w# k
- 如果您的PHP版本>=7.0,可以设置Worker->reusePort=true, 这样可以做到多个子进程创建相同端口的Worker。见下面的例子:
* [' f! R9 g2 c1 `/ ~2 h' Y
2 T1 K2 V* k4 k# W' F2 k- use Workerman\Worker;
' r/ |6 ?0 [0 z0 Y% B6 i# L4 x0 q - require_once './Workerman/Autoloader.php';
% w7 ]7 V0 r3 z8 {" [/ R. F) g2 E; J - 9 _9 B# F( U, B* |% ^% N) W/ }
- $worker = new Worker('text://0.0.0.0:2015');
1 `5 c3 H& y. x9 {* u1 q - // 4个进程
5 S8 }# q, ~0 [* c - $worker->count = 4;. r( v% W5 q' c0 B* |
- // 每个进程启动后在当前进程新增一个Worker监听8 g& W3 K9 M Q2 ~
- $worker->onWorkerStart = function($worker)+ m3 x, N" k& K9 _+ `% I+ P
- {/ R+ C% u$ }8 L/ d% Y% q" ?. }
- $inner_worker = new Worker('http://0.0.0.0:2016');
0 }0 u, D7 B/ N1 }5 g( L! B b - // 设置端口复用,可以创建监听相同端口的Worker(需要PHP>=7.0)
6 g+ X7 h+ Z e7 s( p - $inner_worker->reusePort = true;( \5 v1 ? F8 ?' a4 r
- $inner_worker->onMessage = 'on_message';3 X9 X- r0 |* d; z" G
- // 执行监听。正常监听不会报错
" e }/ v- @4 E6 m - $inner_worker->listen();/ {6 E( {2 }; V3 F
- };3 s G( h4 Q* n1 S+ n- s; C
& U# m# O" W8 R" p- $worker->onMessage = 'on_message';
9 |% _$ t$ z& F! O' E
8 e: u/ j/ M& p* y- function on_message($connection, $data)" ^: U7 S: \9 s0 y" l) q1 C1 B b
- {
7 r: O4 C2 U4 E/ v - $connection->send("hello\n");+ O! B; P" c t+ M4 p/ G
- }
* k7 O: Y& j& N( H4 k" M8 A& b - p, ~" m* J8 ^9 v! s1 Y h; X5 J
- // 运行worker% I/ C) r7 [8 @( Y u
- 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
" L% P9 W6 V3 h: v7 M - use Workerman\Worker;
+ w; c, i% w u+ E' E( ]6 s e - require_once './Workerman/Autoloader.php';
, K; |9 O* c0 g% J - // 初始化一个worker容器,监听1234端口2 R' N' e5 m8 H
- $worker = new Worker('websocket://0.0.0.0:1234');* e* i5 y$ D7 N4 Q, p8 h3 t6 X
7 X3 K3 [) d/ c) b/ V7 H h- /*+ o0 d+ L, I( Z7 I( d; g9 p
- * 注意这里进程数必须设置为1,否则会报端口占用错误" s: H. F4 n M; E+ b6 n
- * (php 7可以设置进程数大于1,前提是$inner_text_worker->reusePort=true)% _5 O+ E% ~6 N0 y; g
- */0 G) A8 g$ z* s
- $worker->count = 1;" b& G$ U: J" a9 Q/ v. ?8 R
- // worker进程启动后创建一个text Worker以便打开一个内部通讯端口
0 F& l0 g2 D1 U" y - $worker->onWorkerStart = function($worker)1 X0 q" B+ G$ B5 m! H+ l$ x
- {
7 Z4 d% t- T8 ?0 O$ {* d - // 开启一个内部端口,方便内部系统推送数据,Text协议格式 文本+换行符
' w" G( H+ D' w. v) a( V8 T - $inner_text_worker = new Worker('text://0.0.0.0:5678');
1 G' L; b2 u! I - $inner_text_worker->onMessage = function($connection, $buffer), H6 Q% [7 V( Z0 I5 q0 [
- {
' _/ c; _/ g4 O. Y; S- V - // $data数组格式,里面有uid,表示向那个uid的页面推送数据% a. a/ m# Z3 i5 f
- $data = json_decode($buffer, true);
* K# m% F& w/ T; \' [0 l6 @. \7 R - $uid = $data['uid'];
) Q' U0 S4 v4 `8 E, s4 M - // 通过workerman,向uid的页面推送数据: R& Y) s0 c' x9 n! x, J6 D$ s
- $ret = sendMessageByUid($uid, $buffer);
G; k& P. ^" T- Y M( _3 @ - // 返回推送结果8 O& T7 B7 v3 I3 R! ?9 l$ J
- $connection->send($ret ? 'ok' : 'fail');
# p X, u" v2 E$ m3 O - };9 u# Z! ]9 v1 T8 M% C# H; ?* ~& m* h
- // ## 执行监听 ##+ p/ ~8 U4 ]! f/ q! X8 I0 m, G
- $inner_text_worker->listen();
. @7 e. j% a7 D7 }' k6 S4 x - };6 S S) x# s( M9 g2 K2 }: X2 H1 _7 j
- // 新增加一个属性,用来保存uid到connection的映射
% Z" @- \- I* h) M - $worker->uidConnections = array();
; i. H% ]! c* D4 S) w; M7 j - // 当有客户端发来消息时执行的回调函数* a4 q& H3 a) O0 |( n
- $worker->onMessage = function($connection, $data)( G1 v: c; V: m1 D" f3 z
- {
/ i+ ]7 G2 L4 h' `8 R - global $worker;( M) D3 b# k: s- c) H
- // 判断当前客户端是否已经验证,既是否设置了uid5 K# Z. k& X) A! B$ j1 {. q! t
- if(!isset($connection->uid))! F1 M! K- s( f4 h4 W
- {- t* X: f j" `
- // 没验证的话把第一个包当做uid(这里为了方便演示,没做真正的验证)
) o5 C( ~$ ?" G8 H1 p - $connection->uid = $data;
. Z5 R* N6 G1 b2 l4 | i - /* 保存uid到connection的映射,这样可以方便的通过uid查找connection,
4 O! j+ J7 ^) ^/ X - * 实现针对特定uid推送数据
2 |! G9 p8 F5 Z' Z+ t; n8 I - */
! `/ e$ D& o+ F8 J, w# S8 j6 O - $worker->uidConnections[$connection->uid] = $connection;
* x( S9 | X; H0 d4 ] - return; O1 Q( W g% Q3 t2 g J
- }+ f6 Y% A" U l1 K
- };* @' F, T8 b9 w: |" v+ v
- : R0 p/ H) a0 V0 `% A6 R8 E1 h, P" h
- // 当有客户端连接断开时! c, ?# i, m1 F
- $worker->onClose = function($connection)5 M( l4 D7 D5 O' X; b
- {
4 `# \* z7 A( Q' p* d: A2 S - global $worker;- r( q6 W1 x0 e l" ]
- if(isset($connection->uid))
2 U% C: d: i: S5 c& L8 Y4 E6 E7 q% D - {; g3 j* G' w3 ~( K7 r/ c
- // 连接断开时删除映射
" `' z7 K) Z+ P2 H J - unset($worker->uidConnections[$connection->uid]);
, N- y! S; @6 a' p5 q - }3 p: j' r9 T/ W
- };
8 d( ~" O( t) s) k" S - 3 | J7 y6 @1 v
- // 向所有验证的用户推送数据6 y, |: r0 r9 A7 M0 n# c
- function broadcast($message)
# K* _; u, E8 ~( Z/ y3 }3 h3 I - {
$ `0 x; T- o0 A! H! s; |3 U- e' n3 W - global $worker;
; ~8 ?2 M8 Y6 K, q - foreach($worker->uidConnections as $connection)4 d+ [/ R: X6 g! x3 B
- {
6 H1 r* o( d+ j( p3 E2 q- _ - $connection->send($message);
( e% A- I! t! B& ^ - }
s( A K) N. @# \/ g - }
; @/ q8 o# P' l$ _5 h
. C p& l1 {$ a, c0 J& }- // 针对uid推送数据; s/ o1 v. l+ a7 t7 u
- function sendMessageByUid($uid, $message)( O; k) r5 @& L2 C
- {
' A, ~* C; z0 G: d3 v" k - global $worker;
; _ e" @; D) @, I: i3 v: \ - if(isset($worker->uidConnections[$uid]))
4 Z* t2 P. f- C. ?. F - {5 b- {3 `7 X: o( |4 @( Q
- $connection = $worker->uidConnections[$uid];9 D6 H8 n' ?9 Z/ R" _* I, V) w
- $connection->send($message);
; e& W* f& Y* M4 {) a9 p - return true;
; d: h& p8 e+ z& }! i" F# q# z6 f - }6 B+ J& d& y. G! k
- return false;
3 C1 J2 g4 G" k: G6 v! f2 U0 V - }" I; Y o4 G4 e9 y- L3 C9 A/ d5 n
8 O2 ]$ M/ ^- G+ |5 I9 C. A- // 运行所有的worker1 {8 A) o" U" [1 }( E$ E9 ~
- Worker::runAll();
复制代码启动后端服务 php push.php start -d 前端接收推送的js代码 - var ws = new WebSocket('ws://127.0.0.1:1234');
0 p; Y. b6 v S8 l: Y8 M - ws.onopen = function(){
4 _' Y( ?: B+ W( n7 w) D: x - var uid = 'uid1';
) F" h. Q: s1 B7 P3 E) |, h - ws.send(uid);4 N7 S8 Q, @3 W
- };5 y. A8 [+ v K0 p( I- S5 F( ~
- ws.onmessage = function(e){
3 W# h* g7 N+ c# ` z9 v - alert(e.data);
5 |3 }" F& `4 B! ] - };
复制代码后端推送消息的代码 - // 建立socket连接到内部推送端口
& W3 u& T6 H w- X0 P R5 p8 P - $client = stream_socket_client('tcp://127.0.0.1:5678', $errno, $errmsg, 1);% v* v) \, R* F: T# D
- // 推送的数据,包含uid字段,表示是给这个uid推送* w; F3 |9 f4 o3 f( c, Y- F
- $data = array('uid'=>'uid1', 'percent'=>'88%');5 L6 J6 v8 B( k2 ?
- // 发送数据,注意5678端口是Text协议的端口,Text协议需要在数据末尾加上换行符
9 R1 @9 Y4 k; u - fwrite($client, json_encode($data)."\n");# x8 D |( E" o
- // 读取推送结果, u0 w, o4 `9 F6 O
- echo fread($client, 8192);
复制代码 ; F" K5 _. g& n
( G7 _. r1 v- f! \ |