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

如何使用Ratchet向特定Socket连接发送事件及管理连接

跨进程向Ratchet Socket连接发送消息的完整方案

嘿,针对你的场景——Ratchet服务已经在管理Socket连接,需要基于客户端的哈希码区分连接,还要从外部PHP进程给指定连接发消息——我整理了一套可行的实现方案,咱们一步步来:

1. 先搞定「区分不同连接」的问题

客户端连接时会发送哈希码做身份识别,那咱们就把这个哈希码和对应的Socket连接绑定起来。修改你的MyApp类,新增一个映射数组来存储「哈希码 → 连接实例」的对应关系,同时在客户端发送认证消息时完成绑定,连接断开时清理映射。

2. 解决「外部PHP访问连接列表」的核心问题

Ratchet是常驻内存的进程,外部PHP是独立的临时进程,不能直接共享内存里的连接实例,所以得用中间件来做跨进程通信。这里推荐用Redis的发布/订阅功能,既高效又容易实现:

  • 让Ratchet进程订阅一个Redis频道,监听外部发来的消息指令
  • 外部PHP进程往这个频道发送包含「目标哈希码」和「消息内容」的指令
  • Ratchet收到指令后,通过之前的哈希-连接映射找到对应连接,发送消息

3. 完整代码实现

首先是修改后的Ratchet应用类

use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;
use Predis\Client;

class MyApp implements MessageComponentInterface {
    protected $connections;
    protected $hashToConnection; // 存储哈希码与连接的映射
    protected $redis;

    public function __construct() {
        $this->connections = new \SplObjectStorage;
        $this->hashToConnection = [];
        
        // 初始化Redis客户端,用来做跨进程通信
        $this->redis = new Client([
            'scheme' => 'tcp',
            'host'   => '127.0.0.1',
            'port'   => 6379,
        ]);

        // 启动Redis订阅循环,监听外部消息
        $this->startRedisSubscription();
    }

    private function startRedisSubscription() {
        $pubSub = $this->redis->pubSubLoop();
        $pubSub->subscribe('socket_targeted_messages');
        
        // 注意:这里要放在ReactPHP的EventLoop里异步执行,避免阻塞Ratchet的主循环
        $loop = \React\EventLoop\Factory::create();
        $loop->addReadStream($pubSub->getConnection()->getResource(), function () use ($pubSub) {
            $message = $pubSub->current();
            $pubSub->next();
            
            if ($message->kind === 'message') {
                $payload = json_decode($message->payload, true);
                // 验证消息格式
                if (isset($payload['hash'], $payload['message']) && isset($this->hashToConnection[$payload['hash']])) {
                    $conn = $this->hashToConnection[$payload['hash']];
                    try {
                        $conn->send($payload['message']);
                    } catch (\Exception $e) {
                        // 连接已断开,清理映射
                        $this->cleanupConnection($conn, $payload['hash']);
                    }
                }
            }
        });
        
        $loop->run();
    }

    public function onOpen(ConnectionInterface $conn) {
        $this->connections->attach($conn);
        echo "新连接建立: (ID {$conn->resourceId})\n";
    }

    public function onMessage(ConnectionInterface $conn, $msg) {
        $data = json_decode($msg, true);
        
        // 处理客户端的认证消息(假设第一个消息是哈希码)
        if (isset($data['auth'])) {
            $userHash = $data['auth'];
            
            // 如果同一个哈希码已有连接,先清理旧连接
            if (isset($this->hashToConnection[$userHash])) {
                $oldConn = $this->hashToConnection[$userHash];
                $this->cleanupConnection($oldConn, $userHash);
            }
            
            // 绑定当前连接与哈希码
            $this->hashToConnection[$userHash] = $conn;
            $conn->send(json_encode(['status' => '认证成功']));
            echo "连接 {$conn->resourceId} 已绑定哈希: {$userHash}\n";
        } else {
            // 处理认证后的业务消息
            echo "收到连接 {$conn->resourceId} 的消息: {$msg}\n";
        }
    }

    public function onClose(ConnectionInterface $conn) {
        $this->cleanupConnection($conn);
        echo "连接 {$conn->resourceId} 已断开\n";
    }

    public function onError(ConnectionInterface $conn, \Exception $e) {
        echo "发生错误: {$e->getMessage()}\n";
        $this->cleanupConnection($conn);
        $conn->close();
    }

    // 封装连接清理逻辑,避免重复代码
    private function cleanupConnection(ConnectionInterface $conn, $userHash = null) {
        $this->connections->detach($conn);
        
        if ($userHash) {
            unset($this->hashToConnection[$userHash]);
        } else {
            // 反向查找哈希码并移除映射
            foreach ($this->hashToConnection as $hash => $connection) {
                if ($connection === $conn) {
                    unset($this->hashToConnection[$hash]);
                    break;
                }
            }
        }
    }
}

然后是外部PHP发送消息的代码

<?php
require 'vendor/autoload.php';

use Predis\Client;

// 连接Redis
$redis = new Client([
    'scheme' => 'tcp',
    'host'   => '127.0.0.1',
    'port'   => 6379,
]);

// 指定要发送的目标哈希码和消息内容
$targetHash = 'user_hash_123';
$message = '来自外部PHP进程的问候!';

// 发送消息到Redis频道
$redis->publish('socket_targeted_messages', json_encode([
    'hash' => $targetHash,
    'message' => $message
]));

echo "消息已发送给哈希码 {$targetHash}\n";

4. 注意事项

  • 先安装Redis客户端依赖:composer require predis/predis
  • 确保Redis服务正在运行,并且配置的连接信息正确
  • 客户端需要在连接建立后立即发送认证哈希码,比如用JS的WebSocket:
    const socket = new WebSocket('ws://your-server-ip:8080');
    socket.onopen = () => {
        socket.send(JSON.stringify({auth: '你的哈希码'}));
    };
    
  • 如果需要给多个哈希码发送消息,可以在外部代码里批量构造消息,或者在Ratchet里扩展批量处理逻辑

内容的提问来源于stack exchange,提问作者Mehdi Azizi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:47:39