ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

5 分钟搞懂:PHP 环境下 Redis Stream 的使用方式

5 分钟搞懂:PHP 环境下 Redis Stream 的使用方式 一、Redis Stream 是什么Redis Stream 是 Redis 5.0 版本正式推出的专用消息队列数据结构。 相比 Redis 传统的消息方案核心差异如下对比Pub/Sub发布订阅模式无持久化、无消息积压、消费者下线即丢失消息仅适合广播通知场景对比ListList 本身支持持久化但仅支持简单的先进先出操作无消费者组概念、无消息确认机制、无法回溯消费、多消费者会重复拉取。Stream 原生支持消息持久化、消费者组、消息确认ACK、消息回溯、死信处理、积压管理等队列特性非常适合需要可靠异步消息的业务场景如订单通知、数据同步、异步任务分发等。二、版本要求Redis Stream 相关能力对服务端、客户端扩展推荐版本组合Redis Server ≥ 5.0.0Stream 数据结构及命令自 5.0 版本正式引入低版本完全不支持。phpredis 扩展 ≥ 4.2.04.2.0 首次完整实现 Streams 系列 API5.3.x 及以上版本稳定性、接口完整性最佳生产环境推荐 5.3.7。PHP ≥ 7.2低版本 PHP 无法安装支持 Stream 的 phpredis 版本。扩展阅读更多版本兼容细节可参考PHP开发必踩的5个Redis版本兼容坑你中了几个附版本对应对照表_php redis扩展6.3.0支持的redis服务端版本-CSDN博客三、Redis 基础连接代码以下为基础连接示例生产环境请根据自身场景调整参数?php $redis new Redis(); try { // 第3参数连接超时时间秒生产环境禁止设为0无限等待 $redis-connect(127.0.0.1, 6379, 2); // 有密码的场景开启 // $redis-auth(your_redis_password); // 选择业务数据库生产环境禁止混用0号默认库 $redis-select(1); // 设置读写超时秒防止慢查询阻塞 PHP 进程 $redis-setOption(Redis::OPT_READ_TIMEOUT, 3); } catch (\RedisException $e) { // 生产环境需做降级处理返回默认值、写入本地缓冲队列等 throw new \RuntimeException(Redis 连接失败: . $e-getMessage()); }四、生产者消息写入的实现生产者负责将业务消息写入 Stream 流核心使用xAdd方法。1. 完整写入示例?php // 1. 初始化连接复用上方连接代码 $redis new Redis(); $redis-connect(127.0.0.1, 6379, 2); $redis-select(1); // 2. 定义流名与消息内容 // 流名推荐用冒号分层按业务域命名 $streamKey order:event:pay_success:stream; // 业务消息体一维关联数组会自动序列化为键值对 $orderData [ order_id NO . date(YmdHis) . mt_rand(1000, 9999), user_id 10086, amount 99.00, ]; // 3. 写入 Stream $maxLen 1000; // 流最大消息数量 $isApproximate true; // 是否近似裁剪 // 消息ID传 * 表示由 Redis 自动生成时间戳-序号全局单调递增 $messageId $redis-xAdd($streamKey, *, $orderData, $maxLen, $isApproximate); echo 消息写入成功流名{$streamKey}\n; echo 消息ID{$messageId}\n;# 通过命令获取指定条数消息xrange streamKey - COUNT 5例如 xrange order:event:pay_success:stream - COUNT 52. 核心参数详解消息 ID*推荐使用 Redis 自动生成格式为「毫秒时间戳 - 序号」全局唯一且单调递增无需业务侧自行生成。MAXLEN 消息上限用于控制流的最大长度避免无限增长占用内存。近似裁剪$isApproximate trueRedis 不会精准卡死在设定条数实际数量会略大于设定值例如示例设置$maxLen 1000;实际数量可能1020或1050性能远高于精确裁剪生产环境推荐开启。3. 生产风险与注意事项⚠️ 重要风险MAXLEN 会直接裁剪最早的消息无论消息是否被消费过。如果消费者速度低于生产速度会直接导致业务消息丢失。可靠业务队列禁止依赖 XADD 自动裁剪应配合积压监控消费完成后手动删除或设置消息过期策略。MAXLEN 仅适用于允许丢失旧消息的场景如日志、实时状态推送。消息体只支持一维关联数组自动序列化如果是嵌套数组、复杂对象必须手动转为 JSON 字符串后写入避免跨语言或解析异常。五、消费者组消息消费的实现消费者组Consumer Group是 Stream 的核心特性同一个流可以创建多个消费组每个组独立消费全量消息组内可以有多个消费者消息自动负载均衡每条消息只会分给组内一个消费者。核心概念说明PELPending Entries List待处理条目列表消费者组维度的待确认消息列表消息被消费者领取后就进入 PEL直到被 ACK 确认。ACKAcknowledgement消息确认消费完成后手动确认消息从 PEL 中移除标记为已处理。MKSTREAM创建消费者组时如果流不存在自动创建空流。1. 示例 1用户通知消费组负责订单支付成功后的短信、站内信通知单消费者即可也可扩展多消费者做负载均衡。?php // 1. 初始化连接 $redis new Redis(); $redis-connect(127.0.0.1, 6379, 2); $redis-select(1); // 2. 定义消费组与消费者 $streamKey order:event:pay_success:stream; $groupName notify_group; // 消费组名一个组对应一个业务域 $consumerName notify_consumer_1; // 消费者名组内需唯一 echo 消费者[{$consumerName}]已启动所属消费组[{$groupName}]\n; echo 业务职责用户通知短信 站内信\n\n; // 3. 创建消费组 // 起始ID 0 表示从流的第一条消息开始消费最后一个参数 true 即 MKSTREAM try { $redis-xGroup(CREATE, $streamKey, $groupName, 0, true); echo ✅ 消费组[{$groupName}]创建成功\n; } catch (\RedisException $e) { // BUSYGROUP 表示组已存在正常跳过即可 if (strpos($e-getMessage(), BUSYGROUP) ! false) { echo ℹ️ 消费组[{$groupName}]已存在跳过创建\n; } else { throw $e; } } // 4. 循环消费消息 while (true) { // XREADGROUP按消费组读取消息 // 表示读取从未分配给该组的新消息 // COUNT 1每次读取1条 // BLOCK 2000无消息时阻塞2秒避免空轮询消耗CPU $messages $redis-xReadGroup( $groupName, $consumerName, [$streamKey ], 1, 2000 ); // 无新消息时继续等待 if (empty($messages)) { echo ⏳ 等待新消息...\n; continue; } // 遍历处理消息 foreach ($messages[$streamKey] as $messageId $fields) { try { echo ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n; echo 收到消息 ID{$messageId}\n; echo 用户ID{$fields[user_id]} | 订单号{$fields[order_id]}\n; // 业务处理1发送短信通知 echo [短信通知] 向用户 {$fields[user_id]} 发送支付成功短信\n; // 业务处理2写入站内信 echo [站内信通知] 向用户 {$fields[user_id]} 写入站内信\n; // 5. 确认消息消费完成 $redis-xAck($streamKey, $groupName, [$messageId]); echo ✅ 消息已确认消费\n\n; } catch (\Throwable $e) { // 单条消息处理失败不终止进程记录日志后继续 echo ❌ 消息处理失败 ID:{$messageId} 错误:{$e-getMessage()}\n; } } }注意因为执行消费之前已经生产3条消息所以执行消费者组后消费了3条2. 示例 2数据同步消费组负责订单数据同步到数仓、更新报表与通知组相互独立各自消费全量消息。?php // 1. 初始化连接 $redis new Redis(); $redis-connect(127.0.0.1, 6379, 2); $redis-select(1); $streamKey order:event:pay_success:stream; $groupName sync_group; $consumerName sync_consumer_1; echo 消费者[{$consumerName}]已启动所属消费组[{$groupName}]\n; echo 业务职责数据同步数仓 销售报表\n\n; // 2. 创建消费组 try { $redis-xGroup(CREATE, $streamKey, $groupName, 0, true); echo ✅ 消费组[{$groupName}]创建成功\n; } catch (\RedisException $e) { if (strpos($e-getMessage(), BUSYGROUP) ! false) { echo ℹ️ 消费组[{$groupName}]已存在跳过创建\n; } else { throw $e; } } // 3. 循环消费 while (true) { $messages $redis-xReadGroup( $groupName, $consumerName, [$streamKey ], 1, 2000 ); if (empty($messages)) { echo ⏳ 等待新消息...\n; continue; } foreach ($messages[$streamKey] as $messageId $fields) { try { echo ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n; echo 收到消息 ID{$messageId}\n; echo 订单号{$fields[order_id]} | 用户ID{$fields[user_id]}\n; // 业务处理1同步到数据仓库ODS层 echo [数仓同步] 将订单 {$fields[order_id]} 同步至ODS层\n; // 业务处理2更新销售报表 echo [报表更新] 更新销售日报数据\n; // 确认消费 $redis-xAck($streamKey, $groupName, [$messageId]); echo ✅ 消息已确认消费\n\n; } catch (\Throwable $e) { echo ❌ 消息处理失败 ID:{$messageId} 错误:{$e-getMessage()}\n; } } }3. 消费进程部署说明⚠️ 生产环境重要提醒 消费进程必须以CLI 命令行模式运行不能通过 PHP-FPM Web 请求执行Web 请求有超时限制无法常驻。 生产环境需配合进程托管工具systemd、supervisor实现异常自动重启同时增加信号监听支持优雅退出。六、进阶1. 消息积压监控使用xInfoGroups查看消费组状态重点关注pending待确认消息数和lag积压数//查询指定 Stream 下所有消费组的统计信息 $groups $redis-xInfo(groups, $streamKey); print_r($groups);返回数组里一共有 2 个消费组notify_group、sync_group下面逐个字段说明。字段含义name消费组名称consumers当前消费组内在线消费者数量pendingPEL 待处理消息数量已经被消费者读取但还没有 ACK 确认的消息last-delivered-id消费组最后一条投递出去的消息 IDentries-read消费组累计读取过的消息总数lag消费组滞后量Stream 中还有多少消息这个消费组还没有消费第2个消费组sync_grouplag7说明Stream 里面还有 7 条消息这个消费组还没有消费存在消息积压。重点区分 pending 和 lagpending已经发给消费者但还没 ACK的消息lag是 Redis 实时计算值表示还没有投递给消费组任何消费者留在 Stream 里的存量消息Redis 会对比消费组last-delivered-id和 Stream 最大消息 ID差值即为 laglag 不为 0 不代表故障要看业务如果是异步同步任务短时 lag 属于正常现象持续上涨则说明消费能力不足pending 上涨则是危险信号消息被消费者拿到但没有 ACK进程崩溃 / 逻辑异常会导致消息重复投递。2. 死信与异常重试处理失败的消息会一直留在 PEL 中可通过xPending查看所有待确认消息。对于多次重试失败的消息建议转移到独立的死信流dead letter stream避免阻塞正常消费后续人工排查。3. 消费幂等性Stream 可能出现消息重复投递如消费者崩溃、网络波动业务侧必须基于业务唯一 ID如订单号做幂等校验避免重复处理。4. 场景异常场景忘记 ACK消息长期积压在 PEL 中占用内存且重启后会重复消费。消费者名不唯一组内消费者重名会导致消息分配混乱出现重复消费。依赖 MAXLEN 裁剪未消费的旧消息被直接删除造成业务数据丢失。Web 模式运行消费进程请求超时后进程被终止消费中断。七、核心总结版本前提Redis Server ≥ 5.0、phpredis ≥ 5.3.x、PHP ≥ 7.2 是生产环境的稳妥组合。核心优势轻量无额外运维成本原生支持持久化、消费者组、ACK 机制适合中小规模异步场景。生产者规范使用自动生成消息 ID可靠业务禁止依赖 MAXLEN 自动裁剪复杂消息手动 JSON 序列化。消费者规范按业务域划分消费组组内消费者名唯一处理逻辑加异常捕获消费完成必须 ACK。生产组合CLI 常驻运行 进程托管 积压监控 幂等校验 死信处理。技术进阶没有捷径但有高效方法。 本号专注分享实战干货、避坑指南、性能调优、面试重难点 每一篇都是亲手落地测试帮你少走弯路、快速提升核心竞争力。❤️ 点赞、在看、收藏、关注一键安排关注账户第一时间获取硬核技术干货下期不见不散
返回列表