如何在PHP中结合RabbitMQ向Ratchet连接的特定用户推送数据
问题描述
我们之前用Ratchet结合ZMQ的SUB/PUB模式实现了WebSocket连接与订阅,数据通过ZMQ上下文推送给订阅用户。现在打算把ZMQ换成RabbitMQ,但不知道集成RabbitMQ时该怎么区分不同的WebSocket连接,从而给特定用户发送数据。相关代码如下:
前端JS代码
var conn = new ab.Session( 'ws://localhost:8000' , function() { conn.subscribe(id, function(topic, data) { console.log(topic); console.log(data); }); } );
后端Server代码
require dirname(__DIR__) . '/vendor/autoload.php'; $loop = React\EventLoop\Factory::create(); $pusher = new MyApp\Pusher; // 监听Web服务器通过Ajax请求后的ZeroMQ推送 $context = new React\ZMQ\Context($loop); $pull = $context->getSocket(ZMQ::SOCKET_PULL); $pull->bind('tcp://localhost:5555'); // 绑定127.0.0.1意味着只有本地客户端能连接 $pull->on('message', array($pusher, 'onBlogEntry')); // 为需要实时更新的客户端搭建WebSocket服务器 $webSock = new React\Socket\Server('0.0.0.0:8000', $loop); // 绑定0.0.0.0允许远程连接 $webServer = new Ratchet\Server\IoServer( new Ratchet\Http\HttpServer( new Ratchet\WebSocket\WsServer( new Ratchet\Wamp\WampServer( $pusher ) ) ), $webSock ); $loop->run();
数据推送代码(Send_data)
$dsn = "tcp://localhost:5555"; $socket = new ZMQSocket(new ZMQContext(), ZMQ::SOCKET_PUSH, 'my pusher'); $endpoints = $socket->getEndpoints(); /* 检查Socket是否已连接 */ if (!in_array($dsn, $endpoints['connect'])) { $socket->connect($dsn); } /* 发送数据 */ $socket->send(json_encode($publish_data));
Pusher类代码
public function onBlogEntry($entry) { $publish_data = json_decode($entry, true); $subscriber_list = $publish_data['subscriber']; if (is_array($subscriber_list) || is_object($subscriber_list)) { foreach ($subscriber_list as $subscriber_list_val) { if (!array_key_exists($subscriber_list_val, $this->subscribedTopics)) { return; }else{ $subscriber = $subscriber_list_val; $topic = $this->subscribedTopics[$subscriber]; $topic->broadcast($publish_data); // 将数据重新发送给订阅该类别的所有客户端 } } } }
替换ZMQ为RabbitMQ的实现方案
核心思路
和ZMQ的PULL/PUSH模式逻辑一致,用RabbitMQ作为中间消息队列接收业务系统的推送消息,再转发给Ratchet的WebSocket服务器。关键是保留原有的用户-订阅映射机制,确保能精准定位到特定WebSocket连接。
步骤1:安装RabbitMQ依赖
composer require php-amqplib/php-amqplib
步骤2:改造WebSocket服务器(替换ZMQ监听为RabbitMQ消费)
修改Server.php,移除ZMQ相关代码,改用RabbitMQ消费者接入消息:
require dirname(__DIR__) . '/vendor/autoload.php'; use PhpAmqpLib\Connection\AMQPStreamConnection; use PhpAmqpLib\Message\AMQPMessage; $loop = React\EventLoop\Factory::create(); $pusher = new MyApp\Pusher; // 连接RabbitMQ服务 $connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest'); $channel = $connection->channel(); // 声明持久化队列(确保队列存在,消息不丢失) $channel->queue_declare('ws_push_queue', false, true, false, false); // 消费RabbitMQ消息,触发原有Pusher逻辑 $channel->basic_consume('ws_push_queue', '', false, true, false, false, function(AMQPMessage $msg) use ($pusher) { $pusher->onBlogEntry($msg->body); }); // 将RabbitMQ消费加入React事件循环,避免阻塞WebSocket服务 $loop->addPeriodicTimer(0.001, function() use ($channel) { while ($channel->is_consuming()) { $channel->wait(null, false); } }); // WebSocket服务器部分保持原有逻辑不变 $webSock = new React\Socket\Server('0.0.0.0:8000', $loop); $webServer = new Ratchet\Server\IoServer( new Ratchet\Http\HttpServer( new Ratchet\WebSocket\WsServer( new Ratchet\Wamp\WampServer( $pusher ) ) ), $webSock ); $loop->run();
步骤3:改造数据推送代码(Send_data)
把ZMQ的PUSH替换为RabbitMQ的消息发布:
use PhpAmqpLib\Connection\AMQPStreamConnection; use PhpAmqpLib\Message\AMQPMessage; $connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest'); $channel = $connection->channel(); // 声明和消费端一致的队列 $channel->queue_declare('ws_push_queue', false, true, false, false); // 构造持久化消息 $msg = new AMQPMessage( json_encode($publish_data), ['delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT] ); // 发送消息到队列 $channel->basic_publish($msg, '', 'ws_push_queue'); // 关闭连接 $channel->close(); $connection->close();
步骤4:保留用户订阅映射逻辑
你的Pusher类中$subscribedTopics的映射机制完全可以复用,它是基于用户ID/订阅主题关联WebSocket连接的,和消息中间件无关。只要业务系统推送的publish_data里包含正确的subscriber列表,就能继续精准推送给特定用户。
关键说明
- RabbitMQ的队列模式和ZMQ的PULL/PUSH等价,都是一对一消息传递,适合当前业务系统向WebSocket服务器推送消息的场景。
- 如果需要更复杂的路由(比如按主题分发),可以改用RabbitMQ的Direct/Topic Exchange模式,但当前场景用简单队列足够。
- React事件循环中的定时任务是为了让RabbitMQ消费不阻塞WebSocket服务的事件循环,保证实时性。
内容的提问来源于stack exchange,提问作者Dev
相关产品推荐
相关产品推荐

