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

使用beyondcodes/laravel-websockets向指定用户推送数据的方法

实现指定用户的WebSocket消息推送方案

针对你使用beyondcode/laravel-websockets时遇到的「无法关联用户与活跃连接、无法给特定用户发消息」的问题,结合你现有的WebSocketHandler代码,给出以下落地实现步骤:

一、核心思路

要给特定用户发消息,必须解决两个核心问题:

  1. 建立用户与WebSocket连接的绑定关系:在连接建立时,将用户ID与连接的socketId关联存储
  2. 维护活跃连接列表:连接断开或异常时清理无效关联,确保只向在线用户推送

推荐用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 02:27:05