使用beyondcodes/laravel-websockets向指定用户推送数据的方法
实现指定用户的WebSocket消息推送方案
针对你使用beyondcode/laravel-websockets时遇到的「无法关联用户与活跃连接、无法给特定用户发消息」的问题,结合你现有的WebSocketHandler代码,给出以下落地实现步骤:
一、核心思路
要给特定用户发消息,必须解决两个核心问题:
- 建立用户与WebSocket连接的绑定关系:在连接建立时,将用户ID与连接的
socketId关联存储 - 维护活跃连接列表:连接断开或异常时清理无效关联,确保只向在线用户推送
推荐用Redis存储关联关系(支持多进程/多服务器场景,内存存储仅适用于单进程部署)。
二、修改WebSocketHandler代码
1. 连接建立时绑定用户并存储关联
修改establishConnection方法,在连接建立时获取用户标识(移动端需在连接时传递,比如user_id或登录token),绑定到连接对象并同步到Redis:
protected function establishConnection(ConnectionInterface $connection) { // 1. 获取用户ID(示例从请求参数取user_id,实际可通过token解析用户) $queryParams = QueryParameters::create($connection->httpRequest); $userId = $queryParams->get('user_id'); // (可选)如果用token验证用户: // $token = $queryParams->get('token'); // $user = Auth::guard('api')->setToken($token)->user(); // if (!$user) throw new UnauthorizedException(); // $userId = $user->id; // 2. 绑定用户ID到连接对象 $connection->userId = $userId; // 3. 存储用户与socketId的关联到Redis // 用集合存储用户的所有活跃socketId(支持同一用户多设备在线) Redis::sadd("websocket:user:{$userId}", $connection->socketId); // 反向映射socketId到用户ID,方便后续清理 Redis::set("websocket:socket:{$connection->socketId}", $userId); // 设置30分钟过期,防止异常断开未清理 Redis::expire("websocket:user:{$userId}", 1800); Redis::expire("websocket:socket:{$connection->socketId}", 1800); // --- 原有逻辑保留(修正原代码的参数获取错误)--- $orderid = $queryParams->get('orderid'); $respones = BModel::select('latitude','longitude') ->where('order_id',$orderid) ->orderBy('created_at','DESC') ->first(); $connection->send(json_encode($respones)); DashboardLogger::connection($connection); StatisticsLogger::connection($connection); return $this; }
2. 连接断开时清理关联
修改onClose方法,移除Redis中对应的关联记录:
public function onClose(ConnectionInterface $connection) { $this->channelManager->removeFromAllChannels($connection); DashboardLogger::disconnection($connection); StatisticsLogger::disconnection($connection); // 清理Redis中的连接关联 if (isset($connection->userId)) { // 从用户的socket集合中移除当前socketId Redis::srem("websocket:user:{$connection->userId}", $connection->socketId); // 如果用户无活跃连接,删除空集合 if (Redis::scard("websocket:user:{$connection->userId}") === 0) { Redis::del("websocket:user:{$connection->userId}"); } } // 删除socketId到用户的映射 Redis::del("websocket:socket:{$connection->socketId}"); }
3. 增加心跳机制维护活跃连接
移动端可能出现异常断开(比如杀进程),导致onClose不触发,需要客户端定期发心跳包,服务端收到后刷新Redis过期时间:
修改onMessage方法处理心跳:
public function onMessage(ConnectionInterface $connection, MessageInterface $message) { $body = collect(json_decode($message->getPayload(), true)); $payload = $body->get('payload', []); // 处理心跳包 if (isset($payload['type']) && $payload['type'] === 'heartbeat') { // 刷新Redis过期时间 if (isset($connection->userId)) { Redis::expire("websocket:user:{$connection->userId}", 1800); } Redis::expire("websocket:socket:{$connection->socketId}", 1800); // 返回心跳响应 $message->respond(json_encode(['type' => 'heartbeat_ack'])); return; } // --- 原有逻辑保留 --- $message->respond(); StatisticsLogger::webSocketMessage($connection); }
三、实现给特定用户发送消息的工具类
创建一个工具类,通过用户ID列表获取活跃连接并发送消息:
<?php namespace App\Services; use BeyondCode\LaravelWebSockets\WebSockets\ConnectionManager; use Illuminate\Support\Facades\Redis; class WebSocketMessageSender { protected $connectionManager; public function __construct(ConnectionManager $connectionManager) { $this->connectionManager = $connectionManager; } /** * 给指定用户列表发送消息 * @param array $userIds 目标用户ID数组 * @param array $message 要发送的消息内容 */ public function sendToUsers(array $userIds, array $message) { $payload = json_encode($message); foreach ($userIds as $userId) { // 获取该用户的所有活跃socketId $socketIds = Redis::smembers("websocket:user:{$userId}"); foreach ($socketIds as $socketId) { // 通过socketId获取对应的连接实例 $connection = $this->connectionManager->getConnection($socketId); if ($connection && $connection->isConnected()) { // 发送消息 $connection->send($payload); } else { // 清理无效的socketId Redis::srem("websocket:user:{$userId}", $socketId); } } } } }
四、使用示例
在控制器或事件中调用工具类发送消息:
// 注入WebSocketMessageSender public function pushMessage(Request $request, WebSocketMessageSender $sender) { // 假设已获取目标用户ID列表 $targetUserIds = [1, 3, 5]; // 要发送的消息内容 $message = [ 'type' => 'notification', 'content' => '您有新的订单通知' ]; $sender->sendToUsers($targetUserIds, $message); return response()->json(['status' => 'success']); }
注意事项
- 确保移动端连接WebSocket时携带用户标识(
user_id或token),否则无法绑定用户 - Redis需正常部署,若用单进程部署,也可使用内存存储(比如静态数组),但不推荐生产环境使用
- 心跳间隔建议设置为5-10分钟,避免频繁请求占用资源
内容的提问来源于stack exchange,提问作者Virat Tripathi
相关产品推荐
相关产品推荐

