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

PHP Ratchet与RabbitMQ整合:实现消息接收并WebSocket广播

我之前在做类似的WebSocket+RabbitMQ推送系统时,也踩过这个“两个监听互相阻塞”的坑。核心问题在于Ratchet本身是基于ReactPHP事件循环的,而默认的RabbitMQ消费代码是同步阻塞的,两个逻辑放在同一个进程里就会互相卡住。下面是我验证过的可行解决方案:

核心思路:共用React事件循环

Ratchet依赖ReactPHP的事件循环来处理WebSocket连接,我们只需要把RabbitMQ的消费逻辑也集成到这个事件循环中,让两者异步运行,就能避免阻塞。这里推荐使用基于React的AMQP客户端库(enqueue/amqp-react),它能完美适配React事件循环。

步骤1:安装依赖

先通过Composer安装所需的库:

composer require enqueue/amqp-react enqueue/amqp-lib ratchet/pawl

步骤2:修改你的Chat类,添加广播方法

确保你的Chat组件有一个可以外部调用的广播方法,用来把RabbitMQ收到的消息推送给所有WebSocket客户端:

use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;

class Chat implements MessageComponentInterface {
    protected $clients;

    public function __construct() {
        $this->clients = new \SplObjectStorage;
    }

    public function onOpen(ConnectionInterface $conn) {
        $this->clients->attach($conn);
        echo "New connection! ({$conn->resourceId})\n";
    }

    public function onMessage(ConnectionInterface $from, $msg) {
        // 保留原有的客户端消息处理逻辑(如果需要)
    }

    public function onClose(ConnectionInterface $conn) {
        $this->clients->detach($conn);
        echo "Connection {$conn->resourceId} has disconnected\n";
    }

    public function onError(ConnectionInterface $conn, \Exception $e) {
        echo "An error occurred: {$e->getMessage()}\n";
        $conn->close();
    }

    // 新增:供外部调用的广播方法
    public function broadcast(string $message) {
        foreach ($this->clients as $client) {
            $client->send($message);
        }
    }
}

步骤3:重构server.php,整合两个逻辑

把Ratchet服务器和RabbitMQ消费者绑定到同一个React事件循环中:

use Ratchet\Server\IoServer;
use Ratchet\Http\HttpServer;
use Ratchet\WebSocket\WsServer;
use React\EventLoop\Loop;
use Enqueue\AmqpReact\AmqpConnectionFactory;

// 初始化Chat组件
$chat = new Chat();

// 获取React主事件循环
$loop = Loop::get();

// 创建并绑定Ratchet服务器到事件循环
$webSocketServer = IoServer::factory(
    new HttpServer(
        new WsServer($chat)
    ),
    8080, // WebSocket监听端口
    '0.0.0.0',
    $loop // 传入主事件循环
);

// 初始化RabbitMQ React客户端
$amqpFactory = new AmqpConnectionFactory([
    'host' => 'localhost',
    'port' => 5672,
    'user' => 'guest',
    'pass' => 'guest',
    'vhost' => '/',
]);

$amqpContext = $amqpFactory->createContext();
$targetQueue = $amqpContext->createQueue('your_message_queue'); // 替换成你的队列名
$amqpContext->declareQueue($targetQueue);

$consumer = $amqpContext->createConsumer($targetQueue);

// 在事件循环中添加周期性的RabbitMQ消息检查(非阻塞)
$loop->addPeriodicTimer(0.001, function () use ($consumer, $chat) {
    // 非阻塞获取队列消息
    if ($message = $consumer->receiveNoWait()) {
        $msgContent = $message->getBody();
        echo "Received message from RabbitMQ: {$msgContent}\n";
        
        // 调用Chat组件的广播方法推送给客户端
        $chat->broadcast($msgContent);
        
        // 确认消息已处理,避免RabbitMQ重复投递
        $consumer->acknowledge($message);
    }
});

// 启动事件循环,同时处理WebSocket和RabbitMQ逻辑
echo "Server started. Listening on ws://0.0.0.0:8080\n";
$loop->run();

关键说明

  1. 非阻塞消费:使用receiveNoWait()替代同步的receive(),这样每次事件循环触发时,只会检查队列有没有新消息,不会卡住进程。
  2. 共用事件循环:Ratchet和RabbitMQ消费逻辑都运行在同一个React事件循环中,实现了异步并发处理。
  3. 消息确认:处理完消息后一定要调用acknowledge(),告诉RabbitMQ消息已被正确处理,避免重复投递。

备选方案:子进程通信

如果不想引入额外库,也可以把RabbitMQ消费逻辑放到独立子进程中,通过管道或Socket和主进程(Ratchet)通信。不过这种方式需要自己处理进程间通信,复杂度更高,不如上面的方案简洁。

内容的提问来源于stack exchange,提问作者Perfect Square

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:36:02