尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
PHP8.5配置WebSocket消息队列怎么实现
前言一条很常见的演进路径第一版用轮询前端每秒发一次 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 位是操作码opcode0x1文本、0x2二进制、0x8关闭、0x9/0xA是 ping/pong掩码方面客户端发来的帧一定带掩码、服务端必须解服务端发出的帧一定不能带掩码方向搞反浏览器会立刻断连。另外帧跑在 TCP 流上一次fread()可能只读到半个帧所以必须为每个连接维护缓冲区循环「取出完整帧剩下的留着」。下面的代码保存为ws-server.phpphp ws-server.php即可运行。它监听两个端口9502给浏览器连9503是内部推送通道队列消费者往这里发消息。?php declare(strict_types1); // 最小可用的 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 还是 StreamRedis 提供三种「队列」语义差别很大选错会在生产上出事故结构命令投递语义订阅者不在线时适合场景列表ListLPUSH/BRPOP一条消息被一个消费者取走消息留着任务分发、点对点推送发布订阅PUBLISH/SUBSCRIBE广播给所有订阅者消息直接丢弃多台 WS 服务器之间的实时扇出流StreamXADD/XREADGROUP消费组至少一次投递 ACK消息留着可回溯可靠投递、需要审计和重放最常用的组合是列表或 Stream 负责可靠投递Pub/Sub 负责跨节点扇出。单机部署只用列表就够了。四、队列消费者把消息推进 WebSocket消费者是独立进程唯一职责是「从队列取消息转交给 WebSocket 服务器」。它和业务之间只通过队列通信重启它不会丢消息。?php declare(strict_types1); // 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 常驻Restartalways保证崩溃自愈再加MemoryMax512M在系统层兜一道内存上限# /etc/systemd/system/ws-server.service 的 [Service] 段 Typesimple Userwww-data ExecStart/usr/bin/php /srv/app/bin/ws-server.php Restartalways RestartSec3 MemoryMax512M别让事件循环做重活。在while (true)里同步查库或调外部接口会阻塞所有连接这类操作要丢给队列同时定时给长时间无响应的连接发 ping超时就dropClient()避免半开连接越积越多。常见坑点❌ 业务代码里直接持有 WebSocket 连接对象并fwrite()推消息✅ 业务只往队列投递投递由独立 worker 承担业务进程重启不影响消息❌ 服务端发出的帧也带掩码或者解析时忘了给客户端帧解掩码✅ 客户端→服务端必须带掩码并解开服务端→客户端必须不带掩码❌ 假设一次fread()就是一个完整帧✅ 必须为每个连接维护缓冲区用pullFrames()循环消费TCP 粘包拆包是常态❌ 用 Pub/Sub 承载不可丢的消息消费者重启期间的消息全部消失✅ 需要不丢就用列表LPUSH/BRPOP或 StreamXADD/XREADGROUPXACK❌ 用一次fwrite()就认为整帧发出去了✅ 非阻塞 socket 上fwrite()的返回值是实际写入字节数可能只写出一半大消息要循环写或维护写缓冲否则客户端会收到截断的帧然后断连❌ 用BRPOP阻塞读却沿用默认的读超时✅setOption(Redis::OPT_READ_TIMEOUT, -1)关掉读超时否则会周期性抛read error on connection❌ 两台服务器各维护连接表消息推不到对面机器上的用户✅ 用 Redis Pub/Sub 扇出每台服务器订阅同一频道任意节点产生的消息都能广播给各自持有的连接总结层次技术选择职责连接层stream_socket_serverstream_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。
RELATED

相关推荐

工作流是什么?从AI节点编排到Coze、Dify、n8n、ComfyUI的通用思维框架

工作流是什么?从AI节点编排到Coze、Dify、n8n、ComfyUI的通用思维框架

如果现在让你用一句话解释“工作流”,你会怎么说?我后台收到最多的问题之一就是“工作流是什么呢”,尤其是最近Coze、Dify、n8n、ComfyUI这些工具轮番刷屏,要么是有人晒出“毛坯房拍照就能生成效果图”的扣子工作流,要…

📅 2026/10/6 9:00:03
dead_code_analyzer鸿蒙适配:Flutter HAP包体积与工程净化

dead_code_analyzer鸿蒙适配:Flutter HAP包体积与工程净化

Flutter 应用跑上鸿蒙之后,我盯着构建产物体积发呆:同样的业务代码,HAP 包比 Android 版大了将近 20%。翻开依赖树和 Dart 源码一查,问题根本不在引擎或中间层,而是大量历史迭代遗留下来的无用类、废弃函数和“暂时注释…

📅 2026/10/6 9:00:03
Flutter for OpenHarmony实战:商品详情页迁移避坑指南

Flutter for OpenHarmony实战:商品详情页迁移避坑指南

做移动开发的同行应该都有感受,2024年下半年开始,Flutter for OpenHarmony这个话题出镜率越来越高。我所在团队做一款二手物品置换App,商品详情页从纯安卓实现迁到Flutter OpenHarmony跑通,前后花了三周。这期间踩了组件通信的坑…

📅 2026/10/6 9:00:03
MORE NEWS

更多资讯

📰

反激变压器设计核心:电感量、磁芯与气隙计算全解析

1. 反激变压器设计的第一道关卡:先把需求翻译成电学参数 做反激变压器设计这事,我踩过最大的坑不是公式算错,而是拿到一个“大概的需求”就急着套公式。比如老板说“做一个5V 2A的手机充电器”,如果你就这么开始算,后面…

📰

MIPI CSI2波形分析实战:从眼图到IBIS仿真的信号完整性定位

1. 从一次调不通的摄像头说起:MIPI CSI2波形分析到底在解决什么问题 摄像头模组点亮失败,十有八九最后都会落到同一个问题上:信号完整性。你可能会先怀疑驱动配置、寄存器时序、I2C通信,甚至换了好几颗模组,结果发现换…

📰

达林顿结构深度解析:从原理到选型避坑指南

很多电路方案里,“达林顿结构”四个字一出现,给人的第一印象就是“放大倍数大、驱动能力强”。我早年间做继电器驱动时也这么想,用单片机的IO口直接推一个TIP122,结果理论上算下来绰绰有余,实际上一接负载就傻眼了&…

📰

白皮书拆解与落地:从目录洞察到实操突围的完整指南

2. 白皮书目录拆解:一份“避坑地图”应该长什么样 既然拿到一份白皮书,第一件事不是从第一章开始细读,而是先看目录。目录就是一份“避坑地图”,也是内容团队留给读者的“寻宝路线”。我拆解过几十份不同行业的白皮书,…

📰

白皮书目录策划:从营销钩子到内容蓝图的全套实操指南

在营销圈混久了,看到“您在寻找突围的利剑吗?白皮书目录已为您备好”这种标题,第一反应不是点开,而是会心一笑——这是非常典型的B2B内容获客手法:前半句用“突围”戳中增长焦虑,后半句用“白皮书目录”给出…

📰

Superpowers插件详解:Atom上TypeScript语言服务与代码智能实战

很多人第一次搜“superpowers”这个词,大概率是被这个名字吸引进来的。在2016年前后的前端圈子里,这个词指的是Atom编辑器上一整套插件体系——它试图把IDE级别的代码智能带进一个原本很轻量的编辑器。我当时是它的重度用户,后来也看着它逐渐…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬