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
相关产品推荐
相关产品推荐

