WebSocket 推消息优先级队列

大白话先说清楚

普通弹幕:  "哈哈哈哈哈"        优先级 1 (低)
礼物打赏:  "送了火箭!"        优先级 2 (中)
系统广播:  "服务器维护通知"    优先级 3 (高)

队列里同时有这三条:
先发系统广播 → 再发礼物 → 最后发弹幕

为什么要优先队列?
直播间一秒几千条弹幕,普通消息把队列塞满了,
礼物消息卡在后面半天没发出去,打赏的人当然骂娘。


代码

<?php

// 优先队列:数字越大越先出队
class MsgQueue extends SplPriorityQueue
{
    // 相同优先级按插入顺序排(先进先出)
    public function compare($a, $b): int
    {
        if ($a['priority'] === $b['priority']) {
            return $a['seq'] < $b['seq'] ? 1 : -1;
        }
        return $a['priority'] < $b['priority'] ? -1 : 1;
    }
}

$server = new Swoole\WebSocket\Server('0.0.0.0', 9502);

$queue = new MsgQueue();
$seq   = 0; // 插入序号,保证同优先级有序

// 优先级常量
define('PRI_CHAT',   1); // 普通弹幕
define('PRI_GIFT',   2); // 礼物打赏
define('PRI_SYSTEM', 3); // 系统广播

$server->on('open', function ($server, $req) {
    echo "fd={$req->fd} 进来了\n";
});

$server->on('message', function ($server, $frame) use (&$queue, &$seq) {
    $data = json_decode($frame->data, true);

    // 根据消息类型定优先级
    $pri = match($data['type'] ?? '') {
        'system' => PRI_SYSTEM,
        'gift'   => PRI_GIFT,
        default  => PRI_CHAT,
    };

    // 塞进队列
    $queue->insert([
        'fd'       => $frame->fd,
        'payload'  => $data,
        'priority' => $pri,
        'seq'      => $seq++,
    ], $pri);

    echo "收到消息 type={$data['type']} pri={$pri}\n";
});

$server->on('close', function ($server, $fd) {
    echo "fd={$fd} 走了\n";
});

// 每10ms消费队列,往所有连接推
$server->on('workerStart', function ($server) use (&$queue) {
    Swoole\Timer::tick(10, function () use ($server, &$queue) {
        // 每次最多发50条,防止一次发太多卡住
        $limit = 50;
        while (!$queue->isEmpty() && $limit-- > 0) {
            $msg = $queue->extract();
            foreach ($server->connections as $fd) {
                if (!$server->isEstablished($fd)) continue;
                $server->push($fd, json_encode($msg['payload']));
            }
        }
    });
});

$server->start();

测试客户端

<?php
// 模拟同时发三种消息

go(function () {
    $c = new Swoole\Coroutine\Http\Client('127.0.0.1', 9502);
    $c->upgrade('/');

    // 同时扔进去,看谁先出来
    $c->push(json_encode(['type' => 'chat',   'msg' => '哈哈哈哈']));
    $c->push(json_encode(['type' => 'gift',   'msg' => '送了火箭']));
    $c->push(json_encode(['type' => 'system', 'msg' => '系统维护']));

    // 接收推过来的消息
    while (true) {
        $frame = $c->recv();
        if ($frame) echo "收到:{$frame->data}\n";
    }
});

队列出队顺序示意

塞入顺序:          出队顺序:
① chat   pri=1     ③ system pri=3  ← 先出
② gift   pri=2     ② gift   pri=2
③ system pri=3     ① chat   pri=1  ← 最后出

核心就三件事

1. SplPriorityQueue   PHP自带优先队列,不用装任何包
2. match 判类型       决定这条消息值几分
3. Timer::tick(10ms)  每10毫秒消费一批,高分的先出去
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐