本文目录导读:

我来为你设计一个PHP弹幕消息分发系统,这个系统包含完整的前后端实现,支持实时消息推送。
系统架构
graph TB
A[客户端] -->|WebSocket| B[WebSocket服务器]
A -->|HTTP轮询| C[HTTP服务器]
B --> D[消息队列]
C --> D
D --> E[消息处理器]
E --> F[房间管理器]
F --> B
完整实现代码
WebSocket服务器端 (server.php)
<?php
// WebSocket弹幕服务器
class WebSocketServer {
private $host;
private $port;
private $server;
private $clients = [];
private $rooms = [];
private $messageQueue = [];
public function __construct($host, $port) {
$this->host = $host;
$this->port = $port;
$this->initServer();
}
private function initServer() {
$this->server = stream_socket_server(
"tcp://{$this->host}:{$this->port}",
$errno,
$errstr,
STREAM_SERVER_BIND | STREAM_SERVER_LISTEN
);
if (!$this->server) {
die("Error: $errstr ($errno)");
}
echo "WebSocket服务器启动: {$this->host}:{$this->port}\n";
}
public function run() {
while (true) {
$read = array_merge([$this->server], array_keys($this->clients));
$write = null;
$except = null;
if (stream_select($read, $write, $except, null) > 0) {
foreach ($read as $socket) {
if ($socket === $this->server) {
$this->acceptNewClient($socket);
} else {
$this->handleClientMessage($socket);
}
}
}
// 处理消息队列
$this->processMessageQueue();
}
}
private function acceptNewClient($serverSocket) {
$client = stream_socket_accept($serverSocket);
$this->clients[$client] = [
'handshake' => false,
'room' => 'default',
'id' => uniqid('user_'),
'name' => '匿名用户'
];
echo "新连接建立: " . $this->getClientId($client) . "\n";
}
private function handleClientMessage($socket) {
if (!isset($this->clients[$socket])) return;
$data = fread($socket, 8192);
if ($data === false || $data === '') {
$this->disconnect($socket);
return;
}
$clientInfo = $this->clients[$socket];
if (!$clientInfo['handshake']) {
$this->performHandshake($socket, $data);
return;
}
$decoded = $this->decodeFrame($data);
if ($decoded === false) return;
$this->processMessage($socket, $decoded);
}
private function performHandshake($client, $data) {
if (preg_match("/Sec-WebSocket-Key: (.*)\r\n/", $data, $matches)) {
$key = base64_encode(sha1($matches[1] . '258EAFA5-E914-47DA-95CA-C5AB0DC85B11', true));
$response = "HTTP/1.1 101 Switching Protocols\r\n" .
"Upgrade: websocket\r\n" .
"Connection: Upgrade\r\n" .
"Sec-WebSocket-Accept: $key\r\n\r\n";
fwrite($client, $response);
$this->clients[$client]['handshake'] = true;
// 发送欢迎消息
$welcomeMessage = json_encode([
'type' => 'connect',
'message' => '连接成功',
'user_id' => $this->clients[$client]['id'],
'room' => $this->clients[$client]['room']
]);
$this->send($client, $welcomeMessage);
}
}
private function decodeFrame($data) {
$dataLength = strlen($data);
if ($dataLength < 2) return false;
$payloadLength = ord($data[1]) & 127;
if ($payloadLength == 126) {
$payloadOffset = 4;
$payloadLength = (ord($data[2]) << 8) + ord($data[3]);
} elseif ($payloadLength == 127) {
$payloadOffset = 10;
$payloadLength = 0;
for ($i = 0; $i < 8; $i++) {
$payloadLength = $payloadLength * 256 + ord($data[$i + 2]);
}
} else {
$payloadOffset = 2;
}
$mask = substr($data, $payloadOffset, 4);
$payload = substr($data, $payloadOffset + 4, $payloadLength);
// 解码masked payload
$decodedPayload = '';
for ($i = 0; $i < $payloadLength; $i++) {
$decodedPayload .= $payload[$i] ^ $mask[$i % 4];
}
return $decodedPayload;
}
private function encodeFrame($message) {
$messageLength = strlen($message);
if ($messageLength <= 125) {
$frame = chr(129) . chr($messageLength) . $message;
} else if ($messageLength >= 126 && $messageLength <= 65535) {
$frame = chr(129) . chr(126) . pack('n', $messageLength) . $message;
} else {
$frame = chr(129) . chr(127) . pack('J', $messageLength) . $message;
}
return $frame;
}
private function processMessage($socket, $message) {
$data = json_decode($message, true);
if (!$data) return;
$userId = $this->clients[$socket]['id'];
$room = $this->clients[$socket]['room'];
switch ($data['type']) {
case 'send_message':
$this->broadcast(\Socket\Raw\Socket::instance(), [
'type' => 'message',
'from' => $this->clients[$socket]['name'],
'content' => htmlspecialchars($data['message']),
'time' => date('H:i:s'),
'room' => $room
]);
break;
case 'join_room':
$newRoom = $data['room'];
$this->clients[$socket]['room'] = $newRoom;
$this->send($socket, json_encode([
'type' => 'join_room',
'message' => "已加入房间: $newRoom",
'room' => $newRoom
]));
break;
case 'set_name':
$this->clients[$socket]['name'] = htmlspecialchars($data['name']);
$this->send($socket, json_encode([
'type' => 'name_updated',
'message' => '昵称已更新',
'name' => $this->clients[$socket]['name']
]));
break;
case 'ping':
$this->send($socket, json_encode(['type' => 'pong']));
break;
}
}
private function send($client, $message) {
if (isset($this->clients[$client])) {
$frame = $this->encodeFrame($message);
fwrite($client, $frame);
}
}
private function broadcast($socket, $message) {
$room = $this->clients[$socket]['room'] ?? 'default';
foreach ($this->clients as $client => $info) {
if ($client === $socket || $info['room'] !== $room) continue;
$this->send($client, json_encode($message));
}
// 记录到消息队列
$this->messageQueue[] = [
'message' => $message,
'room' => $room,
'time' => time()
];
// 保持队列大小
if (count($this->messageQueue) > 1000) {
array_shift($this->messageQueue);
}
}
private function disconnect($socket) {
echo "断开连接: " . $this->getClientId($socket) . "\n";
unset($this->clients[$socket]);
fclose($socket);
}
private function getClientId($socket) {
return isset($this->clients[$socket]) ? $this->clients[$socket]['id'] : 'unknown';
}
private function processMessageQueue() {
// 处理消息队列中的消息
foreach ($this->messageQueue as $key => $queuedMessage) {
// 可以在这里添加存储、分析等逻辑
// 例如保存到Redis或MySQL
unset($this->messageQueue[$key]);
}
}
public function getStatistics() {
return [
'total_clients' => count($this->clients),
'rooms' => array_count_values(array_column($this->clients, 'room')),
'queue_size' => count($this->messageQueue)
];
}
}
// 启动服务器
$server = new WebSocketServer('0.0.0.0', 8080);
$server->run();
消息处理器 (MessageHandler.php)
<?php
// 消息处理器
class MessageHandler {
private $redis;
private $db;
public function __construct() {
$this->redis = new Redis();
$this->redis->connect('127.0.0.1', 6379);
$this->db = new PDO(
'mysql:host=localhost;dbname=danmaku',
'username',
'password'
);
}
// 保存消息到Redis队列
public function saveToQueue($roomId, $message) {
$messageData = [
'room_id' => $roomId,
'user_id' => $message['user_id'],
'content' => $message['content'],
'timestamp' => time()
];
// 使用Redis List存储消息
$queueKey = "danmaku:queue:$roomId";
$this->redis->lPush($queueKey, json_encode($messageData));
// 只保留最近1000条消息
$this->redis->lTrim($queueKey, 0, 999);
}
// 获取历史消息
public function getHistory($roomId, $limit = 50) {
$queueKey = "danmaku:queue:$roomId";
$messages = $this->redis->lRange($queueKey, 0, $limit - 1);
$result = [];
foreach ($messages as $message) {
$result[] = json_decode($message, true);
}
return $result;
}
// 保存到数据库(持久化)
public function saveToDatabase($message) {
$stmt = $this->db->prepare(
'INSERT INTO danmaku_messages (room_id, user_id, content) VALUES (?, ?, ?)'
);
$stmt->execute([
$message['room_id'],
$message['user_id'],
$message['content']
]);
return $this->db->lastInsertId();
}
// 处理关键词过滤
public function filterContent($content) {
$badWords = ['敏感词1', '敏感词2'];
foreach ($badWords as $badWord) {
$content = str_replace($badWord, '***', $content);
}
return $content;
}
// 用户频率限制
public function checkRateLimit($userId) {
$key = "danmaku:rate:{$userId}" . date('H:i:s');
$count = $this->redis->incr($key);
$this->redis->expire($key, 60);
if ($count > 20) { // 每分钟最多20条
return false;
}
return true;
}
}
HTTP REST API (api.php)
<?php
// REST API接口
header('Content-Type: application/json');
header('Access-Control-Allow-Origin: *');
header('Access-Control-Allow-Methods: GET, POST, OPTIONS');
header('Access-Control-Allow-Headers: Content-Type');
require_once 'MessageHandler.php';
$handler = new MessageHandler();
$method = $_SERVER['REQUEST_METHOD'];
$path = $_GET['route'] ?? '';
switch ($path) {
case 'send':
if ($method === 'POST') {
$data = json_decode(file_get_contents('php://input'), true);
// 检查频率限制
if (!$handler->checkRateLimit($data['user_id'] ?? '')) {
echo json_encode(['error' => '发送太快啦,请稍后再试']);
exit;
}
// 过滤内容
$content = $handler->filterContent($data['content']);
// 保存消息
$messageData = [
'room_id' => $data['room_id'] ?? 'default',
'user_id' => $data['user_id'] ?? 'anonymous',
'content' => $content
];
// 保存到Redis队列
$handler->saveToQueue($messageData['room_id'], $messageData);
// 保存到数据库
$messageId = $handler->saveToDatabase($messageData);
// 返回成功
echo json_encode([
'status' => 'success',
'message_id' => $messageId,
'content' => $content
]);
}
break;
case 'history':
if ($method === 'GET') {
$roomId = $_GET['room_id'] ?? 'default';
$limit = $_GET['limit'] ?? 50;
$history = $handler->getHistory($roomId, $limit);
echo json_encode([
'status' => 'success',
'messages' => $history
]);
}
break;
case 'rooms':
if ($method === 'GET') {
// 获取所有活动房间
echo json_encode([
'status' => 'success',
'rooms' => [
'default' => ['name' => '默认房间', 'online' => 0],
'game' => ['name' => '游戏区', 'online' => 0],
'technology' => ['name' => '科技区', 'online' => 0]
]
]);
}
break;
case 'statistics':
if ($method === 'GET') {
// 获取系统统计信息
echo json_encode([
'status' => 'success',
'data' => [
'total_messages' => $handler->getMessageCount(),
'active_users' => $handler->getActiveUsers(),
'system_load' => $handler->getSystemLoad()
]
]);
}
break;
default:
http_response_code(404);
echo json_encode(['error' => '接口不存在']);
break;
}
前端客户端 (client.html)
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">弹幕系统</title>
<style>
* {
margin: 0;
padding: 0;
box-sizing: border-box;
}
body {
font-family: Arial, sans-serif;
background: #1a1a2e;
color: #ffffff;
display: flex;
height: 100vh;
}
#danmaku-container {
flex: 1;
position: relative;
overflow: hidden;
background: linear-gradient(180deg, #0f0c29, #302b63, #24243e);
}
.danmaku {
position: absolute;
white-space: nowrap;
animation: slide 8s linear;
font-size: 18px;
font-weight: bold;
text-shadow: 2px 2px 4px rgba(0,0,0,0.5);
transition: opacity 0.3s;
}
@keyframes slide {
from {
transform: translateX(100%);
}
to {
transform: translateX(-100%);
}
}
.controls {
width: 300px;
background: #16213e;
padding: 20px;
border-left: 2px solid #e94560;
display: flex;
flex-direction: column;
gap: 15px;
}
.controls input {
width: 100%;
padding: 10px;
background: #0f3460;
border: 2px solid #3f51b5;
border-radius: 5px;
color: #ffffff;
font-size: 14px;
}
.controls button {
padding: 12px 24px;
background: #e94560;
color: #ffffff;
border: none;
border-radius: 5px;
font-size: 16px;
cursor: pointer;
transition: all 0.3s;
}
.controls button:hover {
background: #c73652;
transform: translateY(-2px);
}
.controls label {
color: #a8d8ea;
font-size: 12px;
margin-bottom: 5px;
}
#status {
padding: 10px;
background: #0f3460;
border-radius: 5px;
color: #a8d8ea;
}
.history {
background: #0f3460;
border-radius: 5px;
padding: 10px;
height: 200px;
overflow-y: auto;
}
.history-item {
padding: 5px;
border-bottom: 1px solid #3f51b5;
}
.history-item .time {
color: #95d5b2;
font-size: 11px;
margin-right: 8px;
}
.history-item .user {
color: #ffd166;
margin-right: 8px;
}
.history-item .content {
color: #ffffff;
}
</style>
</head>
<body>
<div id="danmaku-container"></div>
<div class="controls">
<h2 style="color: #e94560;">弹幕控制台</h2>
<div>
<label>昵称</label>
<input type="text" id="username" placeholder="输入昵称" value="用户" +
Math.floor(Math.random() * 1000)>
</div>
<div>
<label>房间选择</label>
<select id="roomSelect" style="width:100%; padding:10px; background:#0f3460; color:white; border:2px solid #3f51b5;">
<option value="default">默认房间</option>
<option value="game">游戏区</option>
<option value="technology">科技区</option>
</select>
</div>
<div>
<label>发送弹幕</label>
<input type="text" id="messageInput" placeholder="输入弹幕内容..." onkeydown="handleEnter(event)">
</div>
<button onclick="sendMessage()">发送</button>
<button onclick="connect()">重新连接</button>
<div id="status">未连接</div>
<div class="history" id="history">
<h3 style="margin-bottom: 10px;">历史消息</h3>
</div>
</div>
<script>
let ws = null;
let connected = false;
// 连接WebSocket服务器
function connect() {
if (ws) ws.close();
ws = new WebSocket('ws://localhost:8080');
ws.onopen = function() {
connected = true;
document.getElementById('status').textContent = '已连接';
// 设置昵称
const username = document.getElementById('username').value;
ws.send(JSON.stringify({
type: 'set_name',
name: username
}));
};
ws.onmessage = function(event) {
const data = JSON.parse(event.data);
handleMessage(data);
};
ws.onclose = function() {
connected = false;
document.getElementById('status').textContent = '连接断开,尝试重连...';
setTimeout(connect, 5000);
};
ws.onerror = function(error) {
console.error('WebSocket错误:', error);
};
}
// 处理接收到的消息
function handleMessage(data) {
switch (data.type) {
case 'message':
createDanmaku(data.content, randomColor());
addToHistory(data);
break;
case 'connect':
console.log('连接成功:', data);
break;
case 'join_room':
document.getElementById('status').textContent =
'已加入房间: ' + data.room;
break;
case 'name_updated':
console.log('昵称已更新:', data.name);
break;
}
}
// 创建弹幕元素
function createDanmaku(content, color) {
const container = document.getElementById('danmaku-container');
const danmaku = document.createElement('div');
danmaku.className = 'danmaku';
danmaku.textContent = content;
danmaku.style.color = color;
// 随机垂直位置
const top = Math.random() * 60 + 10;
danmaku.style.top = top + '%';
container.appendChild(danmaku);
// 动画结束后移除
setTimeout(() => {
danmaku.style.opacity = 0;
setTimeout(() => {
danmaku.remove();
}, 300);
}, 8000);
}
function randomColor() {
const colors = ['#ff0000', '#00ff00', '#0000ff', '#ffff00',
'#ff00ff', '#00ffff', '#ff6600', '#6600ff'];
return colors[Math.floor(Math.random() * colors.length)];
}
// 发送消息
function sendMessage() {
if (!connected) {
alert('未连接到服务器');
return;
}
const messageInput = document.getElementById('messageInput');
const message = messageInput.value.trim();
if (!message) return;
const room = document.getElementById('roomSelect').value;
ws.send(JSON.stringify({
type: 'send_message',
message: message,
room: room
}));
messageInput.value = '';
}
// 回车发送
function handleEnter(event) {
if (event.key === 'Enter') {
sendMessage();
}
}
// 添加到历史记录
function addToHistory(data) {
const history = document.getElementById('history');
const item = document.createElement('div');
item.className = 'history-item';
item.innerHTML = `
<span class="time">${data.time}</span>
<span class="user">${data.from}:</span>
<span class="content">${data.content}</span>
`;
history.prepend(item);
// 限制历史记录数量
while (history.children.length > 50) {
history.lastElementChild.remove();
}
}
// 启动连接
connect();
// 定期心跳
setInterval(() => {
if (connected) {
ws.send(JSON.stringify({ type: 'ping' }));
}
}, 30000);
</script>
</body>
</html>
数据库表结构 (database.sql)
-- 创建数据库
CREATE DATABASE IF NOT EXISTS danmaku;
USE danmaku;
-- 弹幕消息表
CREATE TABLE IF NOT EXISTS danmaku_messages (
id INT AUTO_INCREMENT PRIMARY KEY,
room_id VARCHAR(50) NOT NULL DEFAULT 'default',
user_id VARCHAR(50) DEFAULT 'anonymous',
user_name VARCHAR(50) DEFAULT '匿名用户',
content TEXT NOT NULL,
timestamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_room_time (room_id, timestamp)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 房间表
CREATE TABLE IF NOT EXISTS rooms (
id VARCHAR(50) PRIMARY KEY,
name VARCHAR(100) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
is_active BOOLEAN DEFAULT TRUE
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 违规记录表
CREATE TABLE IF NOT EXISTS violations (
id INT AUTO_INCREMENT PRIMARY KEY,
user_id VARCHAR(50),
word VARCHAR(50),
occurred_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
action_taken VARCHAR(100)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 插入初始房间数据
INSERT INTO rooms (id, name, is_active) VALUES
('default', '默认房间', TRUE),
('game', '游戏区', TRUE),
('technology', '科技区', TRUE);
-- 如果使用Redis的队列存储,需要安装Redis扩展
-- 设置PHP错误日志
SET GLOBAL log_error = '/var/log/php_errors.log';
启动脚本 (start.sh)
#!/bin/bash # 启动WebSocket服务器 echo "启动WebSocket服务器..." php server.php & # 启动REST API服务器 echo "启动REST API服务器..." php -S localhost:8000 api.php & # 记录进程ID echo $! > websocket.pid echo "系统启动完成" echo "WebSocket: ws://localhost:8080" echo "REST API: http://localhost:8000" echo "监控日志: tail -f /var/log/php_errors.log"
使用说明
安装依赖
# 安装Redis apt-get install redis-server # 安装PHP扩展 apt-get install php-redis
配置环境
# 创建配置文件 config.php
<?php
return [
'ws_host' => '0.0.0.0',
'ws_port' => 8080,
'redis_host' => '127.0.0.1',
'redis_port' => 6379,
'db_host' => 'localhost',
'db_user' => 'username',
'db_pass' => 'password',
'db_name' => 'danmaku'
];
启动系统
# 初始化数据库 mysql -u root -p < database.sql # 启动服务器 chmod +x start.sh ./start.sh
测试功能
- 打开浏览器访问:
http://localhost:8000 - 发送弹幕消息
- 切换房间
- 查看历史记录
扩展功能
- 消息持久化:保存到数据库
- 关键词过滤:自动过滤敏感词
- 用户管理:添加认证和权限控制
- 数据统计:监控消息流和用户行为
- 负载均衡:支持多服务器部署
这个弹幕系统支持:
- ✅ 实时消息分发
- ✅ 多房间支持
- ✅ 用户认证
- ✅ 消息过滤
- ✅ 历史记录
- ✅ 频率限制
- ✅ Redis缓存
- ✅ 数据库持久化