本文目录导读:

我来提供一个完整的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集成,包括连接管理、消息处理、房间系统、认证等功能,可根据实际需求进行扩展。