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

PHP Ratchet WebSocket从MySQL读取注册用户发送定向消息求助

基于Ratchet实现买家卖家点对点聊天功能改造方案

核心改造思路

  • 绑定WebSocket连接与MySQL中存储的用户ID,建立用户ID到连接对象的映射关系
  • 调整消息传输格式为JSON,携带收信人ID、消息内容等必要字段
  • 新增MySQL交互逻辑,支持拉取已注册用户列表、存储/拉取历史聊天记录

1. 前端代码改造(index.html)

原有代码仅传递了当前登录用户ID作为连接参数,未指定收信人,修改后如下:

<html>
    <head>
        <style>
            input, button { padding: 10px; margin: 5px 0; }
            .msg-item { margin: 8px 0; padding: 6px; background: #f5f5f5; }
        </style>
    </head>
    <body>
        <!-- 实际场景下收信人ID可从卖家详情页动态获取,这里做示例输入 -->
        <input type="text" id="receiver_id" placeholder="收信人用户ID" autocomplete="off" />
        <input type="text" id="message" placeholder="输入消息内容" autocomplete="off" />
        <button onclick="transmitMessage()">发送</button>
        <div id="ids"></div>
        <script>
            // 连接时携带当前登录用户的ID,示例中为450,实际从登录态获取
            var socket  = new WebSocket('ws://localhost:8080?current_user_id=450');
            var messageInput = document.getElementById('message');
            var receiverInput = document.getElementById('receiver_id');
            var msgContainer = document.getElementById('ids');

            function transmitMessage() {
                const msgContent = messageInput.value.trim();
                const receiverId = receiverInput.value.trim();
                if(!msgContent || !receiverId) return;
                // 消息序列化为JSON格式,携带收信人ID
                socket.send(JSON.stringify({
                    receiver_id: parseInt(receiverId),
                    content: msgContent
                }));
                // 自己的消息先渲染到页面
                const selfMsg = document.createElement('div');
                selfMsg.className = 'msg-item';
                selfMsg.textContent = `我: ${msgContent}`;
                msgContainer.appendChild(selfMsg);
                messageInput.value = '';
            }

            socket.onmessage = function(e) {
                const msgData = JSON.parse(e.data);
                const msgItem = document.createElement('div');
                msgItem.className = 'msg-item';
                msgItem.textContent = `用户${msgData.sender_id}: ${msgData.content}`;
                msgContainer.appendChild(msgItem);
            }
        </script>
    </body>
</html>

2. 后端代码改造(socket.php)

新增逻辑说明:

  • 解析连接参数拿到当前登录用户ID,建立用户ID与连接的映射
  • 解析前端发来的JSON消息,直接投递给指定收信人
  • 新增MySQL操作,支持拉取用户列表、存储聊天记录
namespace MyApp;

use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;

class Socket implements MessageComponentInterface {
    // 存储所有连接对象
    protected $clients;
    // 存储用户ID到连接的映射:键为用户ID,值为连接对象
    protected $userConnections = [];
    // MySQL连接对象
    protected $db;

    public function __construct()
    {
        $this->clients = new \SplObjectStorage;
        // 初始化MySQL连接,替换为你自己的数据库配置
        $this->db = new \mysqli('localhost', '数据库用户名', '数据库密码', '数据库名');
        if ($this->db->connect_error) {
            die("数据库连接失败: " . $this->db->connect_error);
        }
        // 设置字符集
        $this->db->set_charset('utf8mb4');
    }

    public function onOpen(ConnectionInterface $conn) {
        // 解析连接URL中的query参数,拿到当前登录用户ID
        $queryParams = [];
        parse_str(parse_url($conn->httpRequest->getUri(), PHP_URL_QUERY), $queryParams);
        $currentUserId = $queryParams['current_user_id'] ?? null;
        
        if(!$currentUserId) {
            $conn->close();
            return;
        }

        // 存储连接
        $this->clients->attach($conn);
        // 绑定用户ID与连接
        $this->userConnections[$currentUserId] = $conn;
        // 给连接对象挂载用户ID属性,后续方便使用
        $conn->userId = $currentUserId;

        echo "新连接接入,用户ID: {$currentUserId}, 资源ID: {$conn->resourceId}\n";

        // 可选:用户上线时推送未读消息
        $unreadSql = "SELECT sender_id, content FROM chat_messages WHERE receiver_id = ? AND is_read = 0";
        $stmt = $this->db->prepare($unreadSql);
        $stmt->bind_param('i', $currentUserId);
        $stmt->execute();
        $result = $stmt->get_result();
        while($msg = $result->fetch_assoc()) {
            $conn->send(json_encode($msg));
        }
        // 标记未读消息为已读
        $updateSql = "UPDATE chat_messages SET is_read = 1 WHERE receiver_id = ? AND is_read = 0";
        $updateStmt = $this->db->prepare($updateSql);
        $updateStmt->bind_param('i', $currentUserId);
        $updateStmt->execute();
    }

    public function onMessage(ConnectionInterface $from, $msg) {
        $msgData = json_decode($msg, true);
        $receiverId = $msgData['receiver_id'] ?? null;
        $content = $msgData['content'] ?? '';

        if(!$receiverId || empty($content)) return;

        // 存储消息到MySQL
        $insertSql = "INSERT INTO chat_messages (sender_id, receiver_id, content, created_at) VALUES (?, ?, ?, NOW())";
        $stmt = $this->db->prepare($insertSql);
        $stmt->bind_param('iis', $from->userId, $receiverId, $content);
        $stmt->execute();

        // 检查收信人是否在线,在线直接推送
        if(isset($this->userConnections[$receiverId])) {
            $receiverConn = $this->userConnections[$receiverId];
            $receiverConn->send(json_encode([
                'sender_id' => $from->userId,
                'content' => $content
            ]));
        }
    }

    public function onClose(ConnectionInterface $conn) {
        // 连接关闭时删除映射关系
        $this->clients->detach($conn);
        if(isset($conn->userId)) {
            unset($this->userConnections[$conn->userId]);
        }
        echo "连接关闭,用户ID: {$conn->userId}\n";
    }

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

    // 额外方法:从MySQL拉取已注册用户列表,你可以根据业务需求调用
    public function getRegisteredUsers() {
        $sql = "SELECT id, username, avatar FROM users WHERE status = 1";
        $result = $this->db->query($sql);
        return $result->fetch_all(MYSQLI_ASSOC);
    }
}

3. 依赖的MySQL表结构

需要提前创建两张表,示例结构如下:

-- 用户表(如果已有用户表可直接复用,字段对应上即可)
CREATE TABLE `users` (
  `id` int unsigned NOT NULL AUTO_INCREMENT,
  `username` varchar(50) NOT NULL,
  `avatar` varchar(255) DEFAULT NULL,
  `status` tinyint NOT NULL DEFAULT '1' COMMENT '1正常 0禁用',
  PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 聊天消息表
CREATE TABLE `chat_messages` (
  `id` int unsigned NOT NULL AUTO_INCREMENT,
  `sender_id` int NOT NULL,
  `receiver_id` int NOT NULL,
  `content` text NOT NULL,
  `created_at` datetime NOT NULL,
  `is_read` tinyint NOT NULL DEFAULT '0' COMMENT '0未读 1已读',
  PRIMARY KEY (`id`),
  KEY `sender_id` (`sender_id`),
  KEY `receiver_id` (`receiver_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

注意事项

  • Ratchet是常驻内存的进程,修改代码后需要重启Socket服务才能生效
  • 长时间运行可能出现MySQL断线问题,可在执行数据库操作前加断线重连逻辑,判断$this->db->ping()是否为true,false的话重新创建连接
  • 如果支持用户多端登录,可以把$userConnections里的结构改成数组,同一个用户ID存储多个连接,发消息时遍历所有连接推送
  • 生产环境建议加消息频率限制,防止恶意刷屏

内容的提问来源于stack exchange,提问作者Salman Akbar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 14:24:02