You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 18:53:16