ThinkPHP项目WebSocket集成方案

wen PHP项目 3

本文目录导读:

ThinkPHP项目WebSocket集成方案

  1. 架构设计
  2. 完整实现代码
  3. 使用方法
  4. 高级功能
  5. 部署注意事项
  6. 安全性配置

我来提供一个完整的ThinkPHP项目WebSocket集成方案。

架构设计

技术选型

  • WebSocket服务器: Workerman (推荐) 或 Swoole
  • TP版本: ThinkPHP 6.0+
  • 通信协议: WebSocket + TCP

目录结构

project/
├── app/
│   ├── websocket/
│   │   ├── controller/
│   │   ├── service/
│   │   └── config/
├── config/
│   └── websocket.php
├── server/
│   └── websocket.php (启动文件)
└── public/
    └── ws/
        └── index.php (前端入口)

完整实现代码

安装依赖

composer require workerman/workerman
composer require workerman/globaldata

配置文件 config/websocket.php

<?php
return [
    // 监听地址
    'listen_ip' => '0.0.0.0',
    'listen_port' => 2346,
    // 心跳配置
    'heartbeat_interval' => 55,
    'heartbeat_timeout' => 120,
    // 进程配置
    'process_count' => 4,
    // 是否开启SSL
    'ssl' => false,
    'ssl_cert' => '',
    'ssl_key' => '',
    // 业务处理器
    'handler' => \app\websocket\service\WebSocketService::class,
    // 鉴权配置
    'auth_enable' => true,
    'auth_class' => \app\websocket\service\JwtAuth::class,
];

WebSocket服务类 app/websocket/service/WebSocketService.php

<?php
namespace app\websocket\service;
use Workerman\Connection\TcpConnection;
use Workerman\Worker;
use app\websocket\service\MessageService;
use app\websocket\service\RoomService;
use app\websocket\service\SessionService;
class WebSocketService
{
    protected $worker;
    protected $messageService;
    protected $roomService;
    protected $sessionService;
    // 连接池
    protected static $connections = [];
    public function __construct()
    {
        $this->messageService = new MessageService();
        $this->roomService = new RoomService();
        $this->sessionService = new SessionService();
    }
    /**
     * 启动WebSocket服务
     */
    public function run()
    {
        $config = config('websocket.');
        // 创建Worker监听端口
        $this->worker = new Worker(
            ($config['ssl'] ? 'websocket+ssl' : 'websocket') . "://{$config['listen_ip']}:{$config['listen_port']}"
        );
        // 设置进程数
        $this->worker->count = $config['process_count'];
        // 设置名称
        $this->worker->name = 'ThinkPHP WebSocket';
        // 设置SSL
        if ($config['ssl']) {
            $this->worker->transport = 'ssl';
            $this->worker->context = [
                'ssl' => [
                    'local_cert' => $config['ssl_cert'],
                    'local_pk' => $config['ssl_key'],
                    'verify_peer' => false
                ]
            ];
        }
        // 连接建立
        $this->worker->onConnect = function ($connection) {
            echo "New connection: {$connection->id}\n";
        };
        // 连接成功
        $this->worker->onWebSocketConnect = function ($connection, $http_header) {
            $this->handleConnection($connection, $http_header);
        };
        // 接收消息
        $this->worker->onMessage = function ($connection, $data) {
            $this->handleMessage($connection, $data);
        };
        // 连接关闭
        $this->worker->onClose = function ($connection) {
            $this->handleClose($connection);
        };
        // 连接错误
        $this->worker->onError = function ($connection, $code, $msg) {
            echo "error $code $msg\n";
        };
        // 进程启动
        $this->worker->onWorkerStart = function ($worker) {
            $this->onWorkerStart($worker);
        };
        // 运行worker
        Worker::runAll();
    }
    /**
     * 处理新连接
     */
    protected function handleConnection($connection, $http_header)
    {
        // 解析token
        $token = $this->parseToken($_GET);
        // 验证身份
        if (config('websocket.auth_enable') && !$this->validateAuth($token)) {
            $connection->close('Authentication failed');
            return;
        }
        // 绑定用户信息
        $connection->uid = $this->getUidFromToken($token);
        $connection->token = $token;
        $connection->last_heartbeat = time();
        // 存储连接
        self::$connections[$connection->id] = $connection;
        // 加入默认房间
        $this->roomService->joinRoom($connection, 'all');
        // 发送欢迎消息
        $this->send($connection, [
            'type' => 'system',
            'message' => 'Welcome to WebSocket Server',
            'time' => date('Y-m-d H:i:s')
        ]);
        // 用户上线通知
        $this->notifyOnlineUsers($connection);
    }
    /**
     * 处理消息
     */
    protected function handleMessage($connection, $data)
    {
        // 更新心跳时间
        $connection->last_heartbeat = time();
        // 解码数据
        $request = json_decode($data, true);
        if (!$request || !isset($request['type'])) {
            $this->sendError($connection, 'Invalid request format');
            return;
        }
        // 心跳处理
        if ($request['type'] === 'ping') {
            $this->send($connection, ['type' => 'pong']);
            return;
        }
        // 根据消息类型分发
        switch ($request['type']) {
            case 'chat':
                $this->messageService->handleChat($connection, $request);
                break;
            case 'broadcast':
                $this->messageService->handleBroadcast($connection, $request);
                break;
            case 'room':
                $this->messageService->handleRoomMessage($connection, $request);
                break;
            case 'private':
                $this->messageService->handlePrivateMessage($connection, $request);
                break;
            case 'join_room':
                $this->roomService->joinRoom($connection, $request['room_id'] ?? '');
                break;
            case 'leave_room':
                $this->roomService->leaveRoom($connection, $request['room_id'] ?? '');
                break;
            default:
                $this->sendError($connection, 'Unknown message type');
        }
    }
    /**
     * 处理连接关闭
     */
    protected function handleClose($connection)
    {
        // 移除连接
        if (isset(self::$connections[$connection->id])) {
            unset(self::$connections[$connection->id]);
        }
        // 离开所有房间
        $this->roomService->leaveAllRooms($connection);
        // 用户下线通知
        if (isset($connection->uid)) {
            $this->notifyOfflineUsers($connection);
            echo "User {$connection->uid} disconnected\n";
        }
        echo "Connection closed: {$connection->id}\n";
    }
    /**
     * 进程启动时的初始化
     */
    protected function onWorkerStart($worker)
    {
        // 设置心跳检测定时器
        if ($worker->id === 0) {
            \Workerman\Timer::add(config('websocket.heartbeat_interval'), function () {
                $this->checkHeartbeats();
            });
        }
        echo "Worker started\n";
    }
    /**
     * 心跳检测
     */
    protected function checkHeartbeats()
    {
        $timeout = config('websocket.heartbeat_timeout');
        foreach (self::$connections as $connection) {
            if (time() - $connection->last_heartbeat > $timeout) {
                $connection->close('Heartbeat timeout');
            }
        }
    }
    /**
     * 发送消息
     */
    public function send($connection, $data)
    {
        $connection->send(json_encode(array_merge([
            'status' => 'success',
            'time' => time()
        ], $data)));
    }
    /**
     * 发送错误
     */
    public function sendError($connection, $message)
    {
        $connection->send(json_encode([
            'status' => 'error',
            'message' => $message,
            'time' => time()
        ]));
    }
    /**
     * 发送给指定用户
     */
    public function sendToUser($uid, $data)
    {
        foreach (self::$connections as $connection) {
            if ($connection->uid === $uid) {
                $this->send($connection, $data);
            }
        }
    }
    /**
     * 发送给所有用户
     */
    public function sendToAll($data)
    {
        foreach (self::$connections as $connection) {
            $this->send($connection, $data);
        }
    }
    /**
     * 发送给房间内所有用户
     */
    public function sendToRoom($roomId, $data)
    {
        $roomConnections = $this->roomService->getRoomConnections($roomId);
        foreach ($roomConnections as $connection) {
            $this->send($connection, $data);
        }
    }
    /**
     * 解析token
     */
    protected function parseToken($params)
    {
        return $params['token'] ?? '';
    }
    /**
     * 验证身份
     */
    protected function validateAuth($token)
    {
        if (empty($token)) {
            return false;
        }
        try {
            $authClass = config('websocket.auth_class');
            $auth = new $authClass();
            return $auth->validate($token);
        } catch (\Exception $e) {
            return false;
        }
    }
    /**
     * 从token获取用户ID
     */
    protected function getUidFromToken($token)
    {
        $auth = new (config('websocket.auth_class'))();
        return $auth->getUid($token);
    }
    /**
     * 用户上线通知
     */
    protected function notifyOnlineUsers($connection)
    {
        if (isset($connection->uid)) {
            $onlineCount = count(self::$connections);
            $this->sendToAll([
                'type' => 'system',
                'event' => 'user_online',
                'message' => "User {$connection->uid} online",
                'online_count' => $onlineCount
            ]);
        }
    }
    /**
     * 用户下线通知
     */
    protected function notifyOfflineUsers($connection)
    {
        if (isset($connection->uid)) {
            $onlines = [];
            foreach (self::$connections as $conn) {
                if (isset($conn->uid) && $conn->id !== $connection->id) {
                    $onlines[] = $conn->uid;
                }
            }
            $this->sendToAll([
                'type' => 'system',
                'event' => 'user_offline',
                'message' => "User {$connection->uid} offline",
                'online_users' => array_unique($onlines)
            ]);
        }
    }
    /**
     * 获取在线用户列表
     */
    public function getOnlineUsers()
    {
        $onlineUsers = [];
        foreach (self::$connections as $connection) {
            if (isset($connection->uid)) {
                $onlineUsers[] = [
                    'uid' => $connection->uid,
                    'last_heartbeat' => $connection->last_heartbeat
                ];
            }
        }
        return $onlineUsers;
    }
    /**
     * 获取连接数
     */
    public function getConnectionCount()
    {
        return count(self::$connections);
    }
    /**
     * 静态获取连接
     */
    public static function getConnection($id)
    {
        return self::$connections[$id] ?? null;
    }
    /**
     * 静态关闭连接
     */
    public static function closeConnection($id)
    {
        if (isset(self::$connections[$id])) {
            self::$connections[$id]->close();
            unset(self::$connections[$id]);
        }
    }
}

消息处理服务 app/websocket/service/MessageService.php

<?php
namespace app\websocket\service;
use app\websocket\service\WebSocketService;
class MessageService
{
    protected $wsServer;
    public function __construct()
    {
        $this->wsServer = new WebSocketService();
    }
    /**
     * 处理聊天消息
     */
    public function handleChat($connection, $request)
    {
        $message = $this->validateMessage($request);
        if (!$message) {
            $this->wsServer->sendError($connection, 'Invalid message');
            return;
        }
        // 保存消息到数据库
        $messageId = $this->saveMessage($connection->uid, $message);
        // 广播给所有用户(或特定房间)
        $this->wsServer->sendToAll([
            'type' => 'chat',
            'message_id' => $messageId,
            'from' => $connection->uid,
            'content' => $message,
            'time' => date('Y-m-d H:i:s')
        ]);
    }
    /**
     * 处理广播消息
     */
    public function handleBroadcast($connection, $request)
    {
        $message = $request['message'] ?? '';
        if (empty($message)) {
            $this->wsServer->sendError($connection, 'Message cannot be empty');
            return;
        }
        $this->wsServer->sendToAll([
            'type' => 'broadcast',
            'from' => $connection->uid,
            'message' => $message,
            'time' => date('Y-m-d H:i:s')
        ]);
    }
    /**
     * 处理房间消息
     */
    public function handleRoomMessage($connection, $request)
    {
        $roomId = $request['room_id'] ?? '';
        $message = $request['message'] ?? '';
        if (empty($roomId) || empty($message)) {
            $this->wsServer->sendError($connection, 'Invalid parameters');
            return;
        }
        $this->wsServer->sendToRoom($roomId, [
            'type' => 'room_message',
            'room_id' => $roomId,
            'from' => $connection->uid,
            'message' => $message,
            'time' => date('Y-m-d H:i:s')
        ]);
    }
    /**
     * 处理私信
     */
    public function handlePrivateMessage($connection, $request)
    {
        $toUid = $request['to_uid'] ?? '';
        $message = $request['message'] ?? '';
        if (empty($toUid) || empty($message)) {
            $this->wsServer->sendError($connection, 'Invalid parameters');
            return;
        }
        // 发送给指定用户
        $this->wsServer->sendToUser($toUid, [
            'type' => 'private_message',
            'from' => $connection->uid,
            'message' => $message,
            'time' => date('Y-m-d H:i:s')
        ]);
        // 回执给发送者
        $this->wsServer->send($connection, [
            'type' => 'private_message_status',
            'status' => 'sent',
            'to_uid' => $toUid,
            'message' => $message
        ]);
    }
    /**
     * 验证消息
     */
    protected function validateMessage($request)
    {
        $message = $request['message'] ?? '';
        // 过滤HTML标签
        $message = strip_tags($message);
        // 去除空白
        $message = trim($message);
        // 限制长度
        if (mb_strlen($message) > 1000) {
            return false;
        }
        return $message;
    }
    /**
     * 保存消息到数据库
     */
    protected function saveMessage($uid, $message)
    {
        // 使用TP的模型保存消息
        try {
            return \think\facade\Db::name('chat_messages')->insertGetId([
                'user_id' => $uid,
                'content' => $message,
                'create_time' => time()
            ]);
        } catch (\Exception $e) {
            return 0;
        }
    }
}

房间管理服务 app/websocket/service/RoomService.php

<?php
namespace app\websocket\service;
class RoomService
{
    // 房间连接池
    protected static $rooms = [];
    /**
     * 加入房间
     */
    public function joinRoom($connection, $roomId)
    {
        if (!isset(self::$rooms[$roomId])) {
            self::$rooms[$roomId] = [];
        }
        self::$rooms[$roomId][$connection->id] = $connection;
        // 通知房间内用户
        $this->broadcastToRoom($roomId, [
            'type' => 'system',
            'event' => 'join_room',
            'user_id' => $connection->uid ?? '',
            'message' => 'User joined room'
        ]);
        return true;
    }
    /**
     * 离开房间
     */
    public function leaveRoom($connection, $roomId)
    {
        if (isset(self::$rooms[$roomId])) {
            unset(self::$rooms[$roomId][$connection->id]);
            // 通知房间内用户
            $this->broadcastToRoom($roomId, [
                'type' => 'system',
                'event' => 'leave_room',
                'user_id' => $connection->uid ?? '',
                'message' => 'User left room'
            ]);
        }
    }
    /**
     * 离开所有房间
     */
    public function leaveAllRooms($connection)
    {
        foreach (self::$rooms as $roomId => &$connections) {
            if (isset($connections[$connection->id])) {
                unset($connections[$connection->id]);
            }
        }
    }
    /**
     * 获取房间内所有连接
     */
    public function getRoomConnections($roomId)
    {
        return self::$rooms[$roomId] ?? [];
    }
    /**
     * 获取房间内用户列表
     */
    public function getRoomUsers($roomId)
    {
        $users = [];
        foreach (self::$rooms[$roomId] ?? [] as $connection) {
            if (isset($connection->uid)) {
                $users[] = $connection->uid;
            }
        }
        return $users;
    }
    /**
     * 广播到房间
     */
    public function broadcastToRoom($roomId, $message)
    {
        foreach (self::$rooms[$roomId] ?? [] as $connection) {
            $connection->send(json_encode($message));
        }
    }
}

认证服务 app/websocket/service/JwtAuth.php

<?php
namespace app\websocket\service;
use Firebase\JWT\JWT;
use Firebase\JWT\Key;
class JwtAuth
{
    protected $secretKey = 'your-secret-key';
    /**
     * 生成Token
     */
    public function generateToken($uid)
    {
        $payload = [
            'uid' => $uid,
            'iat' => time(),
            'exp' => time() + 3600 * 24 // 24小时后过期
        ];
        return JWT::encode($payload, $this->secretKey, 'HS256');
    }
    /**
     * 验证Token
     */
    public function validate($token)
    {
        try {
            JWT::decode($token, new Key($this->secretKey, 'HS256'));
            return true;
        } catch (\Exception $e) {
            return false;
        }
    }
    /**
     * 获取用户ID
     */
    public function getUid($token)
    {
        try {
            $decoded = JWT::decode($token, new Key($this->secretKey, 'HS256'));
            return $decoded->uid;
        } catch (\Exception $e) {
            return null;
        }
    }
}

启动文件 server/websocket.php

<?php
/**
 * WebSocket服务启动文件
 * 运行: php server/websocket.php start
 */
// 定义应用路径
define('APP_PATH', __DIR__ . '/../app/');
define('ROOT_PATH', __DIR__ . '/../');
// 加载自动加载器
require_once __DIR__ . '/../vendor/autoload.php';
// 加载ThinkPHP框架
$app = new \think\App();
$app->initialize();
// 创建WebSocket服务
$wsService = new \app\websocket\service\WebSocketService();
// 运行服务
$wsService->run();

前端JavaScript客户端 public/ws_client.html

<!DOCTYPE html>
<html>
<head>
    <meta charset="UTF-8">WebSocket Client</title>
    <style>
        #messages {
            height: 300px;
            overflow-y: auto;
            border: 1px solid #ccc;
            margin: 10px 0;
            padding: 10px;
        }
        #status {
            color: green;
            font-weight: bold;
            margin: 10px 0;
        }
        #status.disconnected {
            color: red;
        }
    </style>
</head>
<body>
    <h1>WebSocket 测试客户端</h1>
    <div id="status" class="disconnected">未连接</div>
    <button onclick="connect()">连接</button>
    <button onclick="disconnect()">断开</button>
    <div>
        <input type="text" id="message" placeholder="输入消息" size="40">
        <button onclick="sendMessage()">发送</button>
    </div>
    <div id="messages"></div>
    <script>
        var ws = null;
        function connect() {
            // 获取token(这里需要你的实际认证方式)
            var token = 'your-jwt-token';
            // 连接WebSocket
            ws = new WebSocket('ws://localhost:2346?token=' + token);
            ws.onopen = function() {
                document.getElementById('status').innerHTML = '已连接';
                document.getElementById('status').className = '';
                appendMessage('系统', 'WebSocket连接成功');
                // 发送心跳
                setInterval(function() {
                    if (ws.readyState === WebSocket.OPEN) {
                        ws.send(JSON.stringify({type: 'ping'}));
                    }
                }, 30000);
            };
            ws.onmessage = function(event) {
                var data = JSON.parse(event.data);
                appendMessage(data.from || '服务器', data.message || data.content || data);
            };
            ws.onerror = function(error) {
                console.error('WebSocket错误:', error);
                appendMessage('错误', 'WebSocket发生错误');
            };
            ws.onclose = function() {
                document.getElementById('status').innerHTML = '已断开';
                document.getElementById('status').className = 'disconnected';
                appendMessage('系统', 'WebSocket连接已断开');
            };
        }
        function disconnect() {
            if (ws) {
                ws.close();
            }
        }
        function sendMessage() {
            var message = document.getElementById('message').value;
            if (!message || !ws || ws.readyState !== WebSocket.OPEN) {
                return;
            }
            ws.send(JSON.stringify({
                type: 'chat',
                message: message
            }));
            document.getElementById('message').value = '';
        }
        function appendMessage(from, message) {
            var div = document.createElement('div');
            div.innerHTML = '<strong>' + from + ':</strong> ' + message;
            document.getElementById('messages').appendChild(div);
        }
    </script>
</body>
</html>

使用方法

启动服务

# 启动WebSocket服务
php server/websocket.php start
# 启动为守护进程(后台运行)
php server/websocket.php start -d
# 重启服务
php server/websocket.php restart
# 停止服务
php server/websocket.php stop
# 查看状态
php server/websocket.php status

在ThinkPHP控制器中使用

<?php
namespace app\api\controller;
use app\websocket\service\WebSocketService;
use think\Response;
class WsController extends BaseController
{
    /**
     * 发送消息给指定用户
     */
    public function sendMessage()
    {
        $uid = $this->request->post('uid');
        $message = $this->request->post('message');
        $wsSocket = new WebSocketService();
        $wsSocket->sendToUser($uid, [
            'type' => 'system',
            'event' => 'push_message',
            'message' => $message
        ]);
        return json(['success' => true]);
    }
    /**
     * 广播消息
     */
    public function broadcast()
    {
        $message = $this->request->post('message');
        $wsSocket = new WebSocketService();
        $wsSocket->sendToAll([
            'type' => 'system',
            'event' => 'broadcast',
            'message' => $message
        ]);
        return json(['success' => true]);
    }
    /**
     * 获取在线用户
     */
    public function getOnlineUsers()
    {
        $wsSocket = new WebSocketService();
        $users = $wsSocket->getOnlineUsers();
        return json([
            'success' => true,
            'data' => [
                'users' => $users,
                'count' => count($users)
            ]
        ]);
    }
}

高级功能

多进程通信

// 使用GlobalData实现多进程数据共享
use GlobalData\Client;
$globalData = new Client('127.0.0.1:2207');
// 设置共享数据
$globalData->online_count = 10;
// 获取共享数据
$count = $globalData->online_count;

消息队列集成

// 使用Redis发布订阅
use think\facade\Cache;
// 发布消息
Cache::publish('ws_message', json_encode([
    'to' => 'all',
    'data' => $data
]));
// 订阅消息(在worker进程中)
Cache::subscribe('ws_message', function ($message) {
    // 处理消息
});

部署注意事项

Nginx配置

map $http_upgrade $connection_upgrade {
    default upgrade;
    ''      close;
}
upstream websocket_backend {
    server localhost:2346;
}
server {
    listen 80;
    server_name ws.example.com;
    location /ws/ {
        proxy_pass http://websocket_backend;
        proxy_http_version 1.1;
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection $connection_upgrade;
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
        proxy_read_timeout 3600s;
    }
}

Supervisor配置

[program:websocket]
process_name=%(program_name)s
command=php /www/wwwroot/project/server/websocket.php start
autostart=true
autorestart=true
user=www
redirect_stderr=true
stdout_logfile=/www/wwwroot/project/logs/websocket.log
stdout_logfile_maxbytes=10MB

安全性配置

安全配置示例

// 在WebSocketService中添加安全检查
protected function securityCheck($connection)
{
    // 检查IP白名单
    $allowedIPs = config('websocket.allowed_ips');
    $clientIP = $connection->getRemoteIp();
    if (!empty($allowedIPs) && !in_array($clientIP, $allowedIPs)) {
        $connection->close('IP not allowed');
        return false;
    }
    // 检查请求频率
    $rate = new RateLimiter();
    if (!$rate->check($connection->uid, 10, 60)) {
        $connection->close('Too many requests');
        return false;
    }
    return true;
}

这个方案提供了完整的WebSocket集成,包括连接管理、消息处理、房间系统、认证等功能,可根据实际需求进行扩展。

抱歉,评论功能暂时关闭!