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

Ratchet PHP WebSocket负载均衡场景下跨服务器用户通信问题求助

负载均衡下Ratchet WebSocket跨服务器通信解决方案

问题背景

我们部署了3台运行相同PHP包的服务器,通过负载均衡器分配WebSocket连接:用户A1被分配到服务器A,用户B1被分配到服务器B。每个服务器仅在本地存储用户连接的resource id、远程地址等资源对象,导致A1向B1发送消息时,服务器A无法找到B1的连接资源,通信失败。

核心解决方案

利用Redis的发布/订阅(Pub/Sub)机制实现跨节点消息转发,同时将用户连接信息集中存储在Redis中,让所有服务器能共享用户的连接节点信息。

步骤1:给每个服务器配置唯一节点ID

在每个服务器的.env文件中添加唯一标识,用于区分不同节点:

SOCKET_NODE_ID=node-1 # 其他服务器分别设置为node-2、node-3

步骤2:修改SocketServer类,初始化Redis订阅

在SocketServer的构造方法中加入节点ID初始化,并启动Redis订阅:

<?php

namespace App\Lib\Notify;

use App\Lib\Notify\Messages\SocketMessagesFactory;
use App\Models\Mongo\CacheRequest;
use App\Models\TwoFactorAuth;
use Illuminate\Support\Facades\Redis;
use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;
use React\EventLoop\Loop;
use React\EventLoop\LoopInterface;

class SocketServer implements MessageComponentInterface
{
    use SocketAuth, SocketBind,SocketCLIMessages;

    protected $redis;
    protected $nodeId; // 新增节点ID属性
    public static $clients;
    public static $localClients;


    public function __construct()
    {
        $this->redis = Redis::connection();
        $this->nodeId = env('SOCKET_NODE_ID', 'default-node'); // 读取节点ID

        /* Deleting all users from redis related old session */
        self::getKeysLoop('users:*', function ($user, $index) {
            $this->redis->del($index);
        });

        /* Deleting all confirmation reqeusts */
        self::getKeysLoop('confirm:*', function ($user, $index) {
            $this->redis->del($index);
        });

        TwoFactorAuth::truncate();
        CacheRequest::truncate();

        echo self::datetime()."\033[34mDeleted all cached users from online (redis) \033[0m \n";
        
        // 启动Redis订阅,处理跨节点消息
        $this->setupRedisSubscription();
    }

    // 新增Redis订阅方法
    protected function setupRedisSubscription()
    {
        $redisClient = $this->redis->client();
        // 订阅跨服务器消息频道
        $redisClient->subscribe(['websocket:cross-node-messages'], function ($redis, $channel, $message) {
            $msgData = json_decode($message, true);
            if (!$msgData || !isset($msgData['target_resource_id'], $msgData['content'])) {
                return;
            }

            // 检查目标连接是否在当前节点本地,是则直接发送
            if (isset(self::$clients[$msgData['target_resource_id']])) {
                self::$clients[$msgData['target_resource_id']]->send($msgData['content']);
            }
        });

        // 将Redis订阅加入React EventLoop,避免阻塞WebSocket服务
        Loop::addReadStream($redisClient->getSocket(), function ($socket) use ($redisClient) {
            $redisClient->handleSocketRead();
        });
    }

    // 修改onOpen方法,更新连接信息存储逻辑
    public function onOpen(ConnectionInterface $conn)
    {
        $query = $conn->httpRequest->getUri()->getQuery();
        parse_str($query, $queryArr);
        echo self::datetime(). "Connecting ip: ".$conn->remoteAddress." ... \n";
             
        $auth = self::checkToken(($queryArr['token']) ?? null);

        if ($auth['success']) {
            $userId = $auth['user']->data->userId;
            $userRedisKey = 'users:' . $userId;
            // 构造当前连接的节点信息
            $connInfo = [
                'resource_id' => $conn->resourceId,
                'node_id' => $this->nodeId,
                'remote_address' => $conn->remoteAddress
            ];

            $existingUser = $this->redis->hgetall($userRedisKey);
            if ($existingUser) {
                // 追加连接信息到现有用户记录
                $resources = json_decode($existingUser['resources'] ?? '[]', true);
                $resources[] = $connInfo;
                $this->redis->hset($userRedisKey, 'resources', json_encode($resources));
                echo self::datetime().'Push resource to exist user : (userID: '.$userId.' )' . "\n";
            } else {
                // 新建用户在线记录
                $this->redis->hset($userRedisKey, [
                    'user_id' => $userId,
                    'resources' => json_encode([$connInfo])
                ]);
                echo self::datetime().'Connected user (id: '.$userId.') ' . "\n";
                $conn->send('Welcome from Websocket ! Successfully connection.');
                self::bindUser($conn,$queryArr['token'],$auth);
            }

            self::$clients[$conn->resourceId] = $conn;
        }else{
            $conn->close();
            echo self::datetime()."Error in auth: " . $auth['message'] . "\n";
        }
    }

    // 修改onClose方法,清理Redis中的连接信息
    public function onClose(ConnectionInterface $conn)
    {
        // 从本地绑定信息中获取当前连接对应的用户ID(需实现getUserIdByResourceId方法)
        $userId = self::getUserIdByResourceId($conn->resourceId);
        if ($userId) {
            $userRedisKey = 'users:' . $userId;
            $existingUser = $this->redis->hgetall($userRedisKey);
            if ($existingUser) {
                $resources = json_decode($existingUser['resources'], true);
                // 过滤掉当前关闭的连接
                $updatedResources = array_filter($resources, function ($res) use ($conn) {
                    return $res['resource_id'] !== $conn->resourceId;
                });
                $updatedResources = array_values($updatedResources);
                
                if (empty($updatedResources)) {
                    // 用户无剩余连接,删除在线记录
                    $this->redis->del($userRedisKey);
                } else {
                    $this->redis->hset($userRedisKey, 'resources', json_encode($updatedResources));
                }
            }
        }

        self::controllableLogout($conn);
        unset(self::$clients[$conn->resourceId]);
        unset(self::$localClients[$conn->resourceId]);
        echo self::datetime()."Connection closed for ({$conn->remoteAddress}) " . $conn->resourceId . "\n";
    }

    // 保留原有onMessage、onError方法
    public function onMessage(ConnectionInterface $conn, $msg)
    {
        if(is_string($msg)) {
            $message = json_decode($msg);
            if(json_last_error()===JSON_ERROR_NONE) {
                if(isset($message->type) && is_string($message->type) && method_exists(SocketMessagesFactory::class,$message->type)) {
                    $method = $message->type;
                    SocketMessagesFactory::$method($message,$conn);
                }
            }
        }
    }

    public function onError(ConnectionInterface $conn, \Exception $e)
    {
        echo self::datetime()."Error message: " . (string)$e->getMessage() . "\n";
        $conn->close();
    }
}

步骤3:修改消息发送逻辑(SocketMessagesFactory)

在消息处理工厂类中,添加跨节点消息转发逻辑:

<?php

namespace App\Lib\Notify\Messages;

use App\Lib\Notify\SocketServer;
use Illuminate\Support\Facades\Redis;

class SocketMessagesFactory
{
    // 示例:发送消息方法
    public static function sendMessage($message, $conn)
    {
        $targetUserId = $message->target_user_id;
        $redis = Redis::connection();
        $currentNodeId = env('SOCKET_NODE_ID', 'default-node');
        $userRedisKey = 'users:' . $targetUserId;

        $userData = $redis->hgetall($userRedisKey);
        if (!$userData) {
            $conn->send(json_encode(['status' => 'error', 'message' => '用户未在线']));
            return;
        }

        $resources = json_decode($userData['resources'], true);
        $sendContent = json_encode([
            'type' => 'message',
            'sender_id' => self::getSenderIdFromConn($conn), // 需实现获取发送者ID方法
            'data' => $message->data
        ]);

        foreach ($resources as $res) {
            if ($res['node_id'] === $currentNodeId) {
                // 目标连接在本地,直接发送
                if (isset(SocketServer::$clients[$res['resource_id']])) {
                    SocketServer::$clients[$res['resource_id']]->send($sendContent);
                }
            } else {
                // 目标连接在其他节点,发布到Redis频道
                $redis->publish('websocket:cross-node-messages', json_encode([
                    'target_resource_id' => $res['resource_id'],
                    'content' => $sendContent
                ]));
            }
        }
    }

    // 其他消息类型方法...
}

关键注意事项

  • 所有服务器必须连接同一Redis实例/集群,确保发布订阅和数据共享生效
  • Redis订阅需放入React EventLoop异步处理,避免阻塞WebSocket主服务
  • 节点ID必须全局唯一,建议使用服务器hostname或自定义唯一标识
  • 需实现辅助方法:getUserIdByResourceId(从本地绑定信息中获取用户ID)、getSenderIdFromConn(从连接对象中获取发送者ID)
  • 增加错误处理逻辑,应对JSON解码失败、Redis操作异常等情况

内容的提问来源于stack exchange,提问作者Tural Rzaxanov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 05:10:31