PHP 弹幕消息分发

wen PHP项目 3

本文目录导读:

PHP 弹幕消息分发

  1. 系统架构
  2. 完整实现代码
  3. 使用说明
  4. 扩展功能

我来为你设计一个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
  • 发送弹幕消息
  • 切换房间
  • 查看历史记录

扩展功能

  1. 消息持久化:保存到数据库
  2. 关键词过滤:自动过滤敏感词
  3. 用户管理:添加认证和权限控制
  4. 数据统计:监控消息流和用户行为
  5. 负载均衡:支持多服务器部署

这个弹幕系统支持:

  • ✅ 实时消息分发
  • ✅ 多房间支持
  • ✅ 用户认证
  • ✅ 消息过滤
  • ✅ 历史记录
  • ✅ 频率限制
  • ✅ Redis缓存
  • ✅ 数据库持久化

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