前言
一条很常见的演进路径:第一版用轮询,前端每秒发一次 HTTP 请求问「有没有新消息」,用户一多服务器上全是空转请求。第二版上了 WebSocket,连接长住了,但业务代码直接在订单回调里查找连接、fwrite()推消息。于是新问题来了:
- 业务进程和推送进程绑在一起,推送进程一重启,这段时间的消息全部丢失;
- 想扩到两台服务器,结果 A 机器产生的消息推不到连在 B 机器上的用户;
- 进程里执行了一段同步数据库查询,所有连接一起卡住(单线程事件循环);
- 跑了一周后内存一路上扬,最后被 OOM Killer 干掉。
根因是同一个:连接是有状态的,消息却是瞬时事件,两者被绑在了一起。连接属于某台机器的某个进程,消息不属于任何人,中间缺的那一层就是消息队列。
本文用 PHP 8.5 从零搭一套最小方案:纯 PHP 写一个 WebSocket 广播服务器(除核心 stream 函数外不依赖扩展),再用 Redis 队列把「业务产生消息」和「推送给连接」解耦。最低要求PHP 8.1,队列部分需要ext-redis。
一、先把「产生消息」和「投递消息」拆开
| 方案 | 服务器压力 | 实时性 | 进程重启会丢消息吗 | 能否多机广播 |
|---|---|---|---|---|
| HTTP 轮询 / 长轮询 | 高(大量空请求或连接被占住) | 秒级延迟 | 不涉及 | 不涉及 |
| WebSocket 直连业务逻辑 | 低 | 实时 | 会丢 | 不能 |
| WebSocket + 队列 | 低 | 实时 | 不会(消息在队列里) | 能 |
关键的心智转变是:业务代码不应该知道有谁在线,它只负责把消息放进队列。推给哪台机器的哪个连接,由「订阅关系」决定,而订阅关系存在队列里。队列在这一层承担三件事:解耦发送方与投递方、持久化(投递方挂了消息还在)、广播(一条消息送给所有 WebSocket 节点)。
二、用纯 PHP 写一个 WebSocket 服务器
WebSocket 协议(RFC 6455)的核心只有三点:握手时客户端发来Sec-WebSocket-Key,服务端把它拼上一段固定 GUID、做一次 SHA-1 再 base64,回填到Sec-WebSocket-Accept;帧头第一个字节的低 4 位是操作码(opcode),0x1文本、0x2二进制、0x8关闭、0x9/0xA是 ping/pong;掩码方面客户端发来的帧一定带掩码、服务端必须解,服务端发出的帧一定不能带掩码,方向搞反浏览器会立刻断连。另外帧跑在 TCP 流上,一次fread()可能只读到半个帧,所以必须为每个连接维护缓冲区,循环「取出完整帧,剩下的留着」。
下面的代码保存为ws-server.php,php ws-server.php即可运行。它监听两个端口:9502给浏览器连,9503是内部推送通道,队列消费者往这里发消息。
<?php declare(strict_types=1); // 最小可用的 WebSocket 广播服务器(纯 PHP,无扩展依赖) // 启动: php ws-server.php 测试: new WebSocket('ws://127.0.0.1:9502') const LISTEN_HOST = '127.0.0.1'; const LISTEN_PORT = 9502; // 客户端入口 const CONTROL_PORT = 9503; // 内部推送通道(队列消费者连这里) const WS_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'; const MAX_PAYLOAD = 1048576; // 单帧上限 1MB,防止畸形帧吃光内存 $server = stream_socket_server('tcp://' . LISTEN_HOST . ':' . LISTEN_PORT, $e1, $m1); $control = stream_socket_server('tcp://' . LISTEN_HOST . ':' . CONTROL_PORT, $e2, $m2); if ($server === false || $control === false) { fwrite(STDERR, "启动失败: {$m1} / {$m2}\n"); exit(1); } stream_set_blocking($server, false); stream_set_blocking($control, false); $clients = []; // id => ['sock' => resource, 'handshaked' => bool, 'buf' => string] $nextId = 1; fwrite(STDOUT, sprintf( "客户端入口 ws://%s:%d\n内部通道 tcp://%s:%d\n", LISTEN_HOST, LISTEN_PORT, LISTEN_HOST, CONTROL_PORT )); while (true) { $read = [$server, $control]; foreach ($clients as $c) { $read[] = $c['sock']; } $write = $except = null; if (stream_select($read, $write, $except, 1) === false) { fwrite(STDERR, "stream_select 失败\n"); break; } foreach ($read as $sock) { if ($sock === $server) { $conn = @stream_socket_accept($server, 0); if ($conn !== false) { stream_set_blocking($conn, false); $clients[$nextId] = ['sock' => $conn, 'handshaked' => false, 'buf' => '']; $nextId++; } } elseif ($sock === $control) { handleControl($control, $clients); } else { handleClient($sock, $clients); } } } function handleClient($sock, array &$clients): void { $id = null; foreach ($clients as $cid => $c) { if ($c['sock'] === $sock) { $id = $cid; break; } } if ($id === null) { return; } $data = @fread($sock, 65536); if ($data === false || $data === '') { if (feof($sock)) { dropClient($id, $clients); } return; } $clients[$id]['buf'] .= $data; if (!$clients[$id]['handshaked']) { $ok = doHandshake($clients[$id]); if ($ok === null) { return; // 握手请求还没收全 } if ($ok === false) { dropClient($id, $clients); return; } fwrite(STDOUT, sprintf("客户端 #%d 握手完成\n", $id)); } foreach (pullFrames($clients[$id]['buf']) as [$opcode, $payload]) { if ($opcode === 0x1 || $opcode === 0x2) { broadcast($clients, $payload, $opcode); // 演示:原样广播 } elseif ($opcode === 0x8) { sendFrame($clients[$id]['sock'], '', 0x8); dropClient($id, $clients); return; } elseif ($opcode === 0x9) { sendFrame($clients[$id]['sock'], $payload, 0xA); // ping -> pong } } } function dropClient(int $id, array &$clients): void { if (!isset($clients[$id])) { return; } @fclose($clients[$id]['sock']); unset($clients[$id]); } /** @return bool|null null 表示数据不完整,false 表示握手非法 */ function doHandshake(array &$client): ?bool { $pos = strpos($client['buf'], "\r\n\r\n"); if ($pos === false) { return null; } $header = substr($client['buf'], 0, $pos); $client['buf'] = substr($client['buf'], $pos + 4); if (!preg_match('/^Sec-WebSocket-Key:\s*(.+)$/mi', $header, $m)) { return false; } $accept = base64_encode(sha1(trim($m[1]) . WS_GUID, true)); $response = "HTTP/1.1 101 Switching Protocols\r\n" . "Upgrade: websocket\r\n" . "Connection: Upgrade\r\n" . "Sec-WebSocket-Accept: {$accept}\r\n\r\n"; if (@fwrite($client['sock'], $response) === false) { return false; } $client['handshaked'] = true; return true; } /** 从缓冲区取出所有完整帧,剩余字节留在 $buf 里 */ function pullFrames(string &$buf): array { $frames = []; while (strlen($buf) >= 2) { $opcode = ord($buf[0]) & 0x0F; $masked = (ord($buf[1]) & 0x80) !== 0; $payloadLen = ord($buf[1]) & 0x7F; $offset = 2; if ($payloadLen === 126) { if (strlen($buf) < 4) break; $payloadLen = unpack('n', substr($buf, 2, 2))[1]; $offset = 4; } elseif ($payloadLen === 127) { if (strlen($buf) < 10) break; $payloadLen = unpack('J', substr($buf, 2, 8))[1]; $offset = 10; } if ($payloadLen > MAX_PAYLOAD) { $buf = ''; // 畸形帧,直接丢弃 return $frames; } $maskKey = ''; if ($masked) { if (strlen($buf) < $offset + 4) break; $maskKey = substr($buf, $offset, 4); $offset += 4; } if (strlen($buf) < $offset + $payloadLen) break; // 帧还没收全 $payload = substr($buf, $offset, $payloadLen); if ($masked && $payloadLen > 0) { $payload ^= str_repeat($maskKey, (int)ceil($payloadLen / 4)); $payload = substr($payload, 0, $payloadLen); } $buf = substr($buf, $offset + $payloadLen); $frames[] = [$opcode, $payload]; } return $frames; } /** 服务端发出的帧不带掩码;FIN 置 1,只发单帧 */ function sendFrame($sock, string $payload, int $opcode = 0x1): void { $len = strlen($payload); if ($len < 126) { $head = chr(0x80 | $opcode) . chr($len); } elseif ($len <= 0xFFFF) { $head = chr(0x80 | $opcode) . chr(126) . pack('n', $len); } else { $head = chr(0x80 | $opcode) . chr(127) . pack('J', $len); } // 演示用单次写入;生产环境要处理「只写出去一半」的情况(见坑点 5) @fwrite($sock, $head . $payload); } function broadcast(array $clients, string $payload, int $opcode = 0x1): void { foreach ($clients as $c) { if ($c['handshaked'] && is_resource($c['sock'])) { sendFrame($c['sock'], $payload, $opcode); } } } /** 队列消费者连上 9503、写一条消息、断开;这里读到就广播出去 */ function handleControl($control, array &$clients): void { $conn = @stream_socket_accept($control, 0); if ($conn === false) { return; } stream_set_timeout($conn, 1); // 防止慢连接卡死事件循环 $payload = ''; while (!feof($conn) && strlen($payload) < MAX_PAYLOAD) { $chunk = fread($conn, 65536); if ($chunk === false || $chunk === '') { break; } $payload .= $chunk; } fclose($conn); $payload = trim($payload); if ($payload === '') { return; } fwrite(STDOUT, sprintf("推送: %s\n", $payload)); broadcast($clients, $payload); }三、队列层:列表、Pub/Sub 还是 Stream?
Redis 提供三种「队列」,语义差别很大,选错会在生产上出事故:
| 结构 | 命令 | 投递语义 | 订阅者不在线时 | 适合场景 |
|---|---|---|---|---|
| 列表(List) | LPUSH/BRPOP | 一条消息被一个消费者取走 | 消息留着 | 任务分发、点对点推送 |
| 发布订阅 | PUBLISH/SUBSCRIBE | 广播给所有订阅者 | 消息直接丢弃 | 多台 WS 服务器之间的实时扇出 |
| 流(Stream) | XADD/XREADGROUP | 消费组,至少一次投递 + ACK | 消息留着,可回溯 | 可靠投递、需要审计和重放 |
最常用的组合是:列表或 Stream 负责可靠投递,Pub/Sub 负责跨节点扇出。单机部署只用列表就够了。
四、队列消费者:把消息推进 WebSocket
消费者是独立进程,唯一职责是「从队列取消息,转交给 WebSocket 服务器」。它和业务之间只通过队列通信,重启它不会丢消息。
<?php declare(strict_types=1); // queue-worker.php —— 从 Redis 取消息,推给 WebSocket 服务器 // 启动: php queue-worker.php $redis = new Redis(); $redis->connect('127.0.0.1', 6379, 2.0); // 关键:BRPOP 阻塞等待时不能有读超时,否则会不停抛异常 $redis->setOption(Redis::OPT_READ_TIMEOUT, -1); fwrite(STDOUT, "worker 已启动,等待消息...\n"); while (true) { try { // 阻塞取一条,返回 [队列名, 消息体] $item = $redis->brPop(['ws:broadcast', 'ws:direct'], 5); if (!is_array($item) || count($item) < 2) { continue; // 超时,回到循环顶部 } [$queue, $payload] = $item; $sock = @stream_socket_client('tcp://127.0.0.1:9503', $errno, $errstr, 2.0); if ($sock === false) { // 生产环境:应把消息写回队列或写进重试队列,而不是直接丢掉 fwrite(STDERR, "推送失败: {$errstr}\n"); continue; } fwrite($sock, $payload . "\n"); fclose($sock); } catch (RedisException $e) { fwrite(STDERR, "Redis 异常: {$e->getMessage()}\n"); sleep(1); $redis->connect('127.0.0.1', 6379, 2.0); // 退避后重连 } }业务侧发布消息只有一行,完全不需要知道谁在线:
// 业务侧:支付成功后往队列里丢一条消息 $redis->lPush('ws:broadcast', json_encode( ['type' => 'order.paid', 'orderId' => 10231], JSON_UNESCAPED_UNICODE | JSON_THROW_ON_ERROR ));4.1 长驻进程的运维要点
WebSocket 服务器和队列消费者都是长驻进程,和「跑一次就退出」的脚本有本质区别。三条纪律:
- 显式设置
memory_limit。CLI 下它的默认策略和 FPM 不同(通常是不限制),一旦泄漏就会把整台机器吃光。入口脚本里ini_set('memory_limit', '512M'),并周期性打印memory_get_usage(true)观察曲线。 - 让进程被托管。用 systemd 常驻,
Restart=always保证崩溃自愈,再加MemoryMax=512M在系统层兜一道内存上限:
# /etc/systemd/system/ws-server.service 的 [Service] 段 Type=simple User=www-data ExecStart=/usr/bin/php /srv/app/bin/ws-server.php Restart=always RestartSec=3 MemoryMax=512M- 别让事件循环做重活。在
while (true)里同步查库或调外部接口会阻塞所有连接,这类操作要丢给队列;同时定时给长时间无响应的连接发 ping,超时就dropClient(),避免半开连接越积越多。
常见坑点
- ❌ 业务代码里直接持有 WebSocket 连接对象并
fwrite()推消息
✅ 业务只往队列投递,投递由独立 worker 承担,业务进程重启不影响消息
- ❌ 服务端发出的帧也带掩码,或者解析时忘了给客户端帧解掩码
✅ 客户端→服务端必须带掩码并解开,服务端→客户端必须不带掩码
- ❌ 假设一次
fread()就是一个完整帧
✅ 必须为每个连接维护缓冲区,用pullFrames()循环消费;TCP 粘包拆包是常态
- ❌ 用 Pub/Sub 承载不可丢的消息,消费者重启期间的消息全部消失
✅ 需要不丢就用列表(LPUSH/BRPOP)或 Stream(XADD/XREADGROUP+XACK)
- ❌ 用一次
fwrite()就认为整帧发出去了
✅ 非阻塞 socket 上fwrite()的返回值是实际写入字节数,可能只写出一半;大消息要循环写或维护写缓冲,否则客户端会收到截断的帧然后断连
- ❌ 用
BRPOP阻塞读,却沿用默认的读超时
✅setOption(Redis::OPT_READ_TIMEOUT, -1)关掉读超时,否则会周期性抛read error on connection
- ❌ 两台服务器各维护连接表,消息推不到对面机器上的用户
✅ 用 Redis Pub/Sub 扇出:每台服务器订阅同一频道,任意节点产生的消息都能广播给各自持有的连接
总结
| 层次 | 技术选择 | 职责 |
|---|---|---|
| 连接层 | stream_socket_server+stream_select | 维护连接、握手与帧解析 |
| 传输层 | 内部 TCP 端口(9503) | 接收待广播的消息 |
| 队列层 | Redis List / Stream / Pub-Sub | 解耦、持久化、跨节点扇出 |
| 消费层 | 独立 PHP worker 进程 | 从队列取消息并转交 |
| 业务层 | 一行LPUSH | 只管投递,不知道谁在线 |
| 运维层 | systemd +MemoryMax+ 心跳 | 保活、限内存、清理死连接 |
技术核心不是「怎么写帧」,而是把连接和消息彻底分开:连接归 WebSocket 服务器管,消息归队列管,两者之间靠一个消费者进程连接。守住这条边界,进程重启、水平扩容、消息重放都会自然满足。最后提醒:PHP 8.5 并没有改动 socket 与 stream 的 API,但如果你打算用 Workerman、Swoole、Ratchet 这类成熟框架替代本文的手写实现,上线前先确认框架声明的 PHP 版本支持范围已经覆盖 8.5。