swoole方案 WebSocket 弱网断线重连与消息补偿机制
·
正常:服务器推消息 msg_id=1,2,3,4,5
断网:msg_id=3,4,5 没收到
重连:客户端说"我最后收到2",服务器查Redis补发3,4,5
漫游队列:Redis 里存最近 N 条消息,按 msg_id 排序。客户端带着"水位线"重连,服务器把水位线之后的全补过去。
---
代码
<?php
$users = new Swoole\Table(1 << 16);
$users->column('room', Swoole\Table::TYPE_STRING, 32);
$users->column('uid', Swoole\Table::TYPE_STRING, 64);
$users->create();
$server = new Swoole\WebSocket\Server('0.0.0.0', 9502);
$server->set(['worker_num' => 4, 'dispatch_mode' => 2]);
function redis(): Swoole\Coroutine\Redis {
$r = new Swoole\Coroutine\Redis();
$r->connect('127.0.0.1', 6379);
return $r;
}
// 连接/重连:带上 uid、room、最后收到的 msg_id
$server->on('open', function ($server, $req) use ($users) {
$uid = $req->get['uid'] ?? 'u' . $req->fd;
$room = $req->get['room'] ?? 'hall';
$lastMsgId = (int)($req->get['last'] ?? 0);
$users->set($req->fd, ['room' => $room, 'uid' => $uid]);
// 补发断线期间漏掉的消息(last+1 到最新)
if ($lastMsgId >= 0) {
$r = redis();
$msgs = $r->zRangeByScore("msgs:$room", $lastMsgId + 1, '+inf');
foreach (($msgs ?: []) as $raw) {
$server->push($req->fd, $raw);
}
}
});
// 收到客户端发来的消息,存队列再广播
$server->on('message', function ($server, $frame) use ($users) {
$data = json_decode($frame->data, true);
$row = $users->get($frame->fd);
$room = $row['room'];
$r = redis();
$msgId = $r->incr("msgid:$room"); // 全局自增ID,单调递增
$msg = json_encode([
'msg_id' => $msgId,
'uid' => $row['uid'],
'content' => $data['content'] ?? '',
'time' => time(),
]);
// 存进漫游队列,Sorted Set,score=msgId
$r->zAdd("msgs:$room", $msgId, $msg);
$r->zRemRangeByRank("msgs:$room", 0, -501); // 只保留最新500条,超出的裁掉
// 广播给房间所有在线的人
foreach ($users as $fd => $u) {
if ($u['room'] === $room && $server->isEstablished($fd)) {
$server->push($fd, $msg);
}
}
});
$server->on('close', fn($s, $fd) => $users->del($fd));
$server->start();
---
前端(带重连和水位线)
class RobustWS {
constructor(uid, room) {
this.uid = uid;
this.room = room;
this.lastMsgId = +(localStorage.getItem(`last:${room}`) ?? 0);
this.delay = 1000; // 初始重连间隔
this.connect();
}
connect() {
const url = `ws://localhost:9502?uid=${this.uid}&room=${this.room}&last=${this.lastMsgId}`;
this.ws = new WebSocket(url);
this.ws.onopen = () => {
console.log(`重连成功,水位线: ${this.lastMsgId}`);
this.delay = 1000; // 重连成功,间隔重置
};
this.ws.onmessage = ({ data }) => {
const msg = JSON.parse(data);
// 去重:补发的消息可能和本地已有的重叠
if (msg.msg_id <= this.lastMsgId) return;
this.lastMsgId = msg.msg_id;
localStorage.setItem(`last:${this.room}`, msg.msg_id); // 水位线持久化
this.onMessage(msg);
};
this.ws.onclose = () => {
// 指数退避:1s -> 2s -> 4s -> 最多30s
setTimeout(() => this.connect(), this.delay);
this.delay = Math.min(this.delay * 2, 30000);
};
}
send(content) {
if (this.ws.readyState === 1) this.ws.send(JSON.stringify({ content }));
}
onMessage(msg) {} // 业务层覆盖
}
// 用法
const ws = new RobustWS('user1', 'hall');
ws.onMessage = msg => console.log(`[${msg.msg_id}] ${msg.uid}: ${msg.content}`);
ws.send('hello');
---
解释
msg_id 是水位线:单调递增的序号,记录"我收到多远了"。客户端把它存 localStorage,断电重启也不丢。重连时带上,服务器就知道从哪里补
Redis INCR:全局原子自增,4个 worker 进程同时调用不会冲突,保证 msg_id 不重复不乱序
Sorted Set zAdd(score=msgId):按 msg_id 排序存消息。zRangeByScore(last+1, +inf) 一句话拿出所有漏掉的消息,不用遍历
zRemRangeByRank(0,
-501):只留最新500条,第501条进来就把最老的踢掉。漫游队列不是全量存储,太早的消息补不了,提示用户"消息已过期请刷新"
客户端去重 if msg_id <= lastMsgId:补发的消息区间可能和客户端本地缓存重叠,前端收到就检查,老的扔掉,防止消息重复显示
指数退避:断了1秒重连,还断2秒,还断4秒,最多30秒试一次。不这样的话弱网下几千个客户端同时疯狂重连,服务器直接被打垮
---
测试
php server.php
# 客户端A进房间
wscat -c "ws://localhost:9502?uid=user1&room=hall&last=0"
# 客户端B发几条消息(模拟A断线期间)
wscat -c "ws://localhost:9502?uid=user2&room=hall&last=0"
> {"content":"消息1"}
> {"content":"消息2"}
> {"content":"消息3"}
# A重连,带上最后收到的水位线
wscat -c "ws://localhost:9502?uid=user1&room=hall&last=1"
# 立刻收到补发: msg_id=2的消息, msg_id=3的消息
更多推荐


所有评论(0)