一、Redis Stream 是什么
Redis Stream 是 Redis 5.0 版本正式推出的专用消息队列数据结构。 相比 Redis 传统的消息方案,核心差异如下:
- 对比Pub/Sub:发布订阅模式无持久化、无消息积压、消费者下线即丢失消息,仅适合广播通知场景;
- 对比List:List 本身支持持久化,但仅支持简单的先进先出操作,无消费者组概念、无消息确认机制、无法回溯消费、多消费者会重复拉取。
Stream 原生支持消息持久化、消费者组、消息确认(ACK)、消息回溯、死信处理、积压管理等队列特性,非常适合需要可靠异步消息的业务场景,如订单通知、数据同步、异步任务分发等。
二、版本要求
Redis Stream 相关能力对服务端、客户端扩展,推荐版本组合:
- Redis Server ≥ 5.0.0:Stream 数据结构及命令自 5.0 版本正式引入,低版本完全不支持。
- phpredis 扩展 ≥ 4.2.0:4.2.0 首次完整实现 Streams 系列 API;5.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 = true:Redis 不会精准卡死在设定条数,实际数量会略大于设定值(例如示例设置$maxLen = 1000;实际数量可能1020或1050),性能远高于精确裁剪,生产环境推荐开启。
3. 生产风险与注意事项
⚠️ 重要风险:MAXLEN 会直接裁剪最早的消息,无论消息是否被消费过。如果消费者速度低于生产速度,会直接导致业务消息丢失。
- 可靠业务队列禁止依赖 XADD 自动裁剪,应配合积压监控,消费完成后手动删除或设置消息过期策略。
- MAXLEN 仅适用于允许丢失旧消息的场景,如日志、实时状态推送。
- 消息体只支持一维关联数组自动序列化;如果是嵌套数组、复杂对象,必须手动转为 JSON 字符串后写入,避免跨语言或解析异常。
五、消费者组:消息消费的实现
消费者组(Consumer Group)是 Stream 的核心特性:同一个流可以创建多个消费组,每个组独立消费全量消息;组内可以有多个消费者,消息自动负载均衡,每条消息只会分给组内一个消费者。
核心概念说明
- PEL(Pending Entries List,待处理条目列表):消费者组维度的待确认消息列表,消息被消费者领取后就进入 PEL,直到被 ACK 确认。
- ACK(Acknowledgement,消息确认):消费完成后手动确认,消息从 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 | 当前消费组内在线消费者数量 |
| pending | PEL 待处理消息数量(已经被消费者读取,但还没有 ACK 确认的消息) |
| last-delivered-id | 消费组最后一条投递出去的消息 ID |
| entries-read | 消费组累计读取过的消息总数 |
| lag | 消费组滞后量:Stream 中还有多少消息,这个消费组还没有消费 |
第2个消费组sync_group:lag=7说明Stream 里面还有 7 条消息,这个消费组还没有消费,存在消息积压。
重点区分 pending 和 lag:
- pending:已经发给消费者,但还没 ACK的消息;
lag是 Redis 实时计算值,表示还没有投递给消费组任何消费者,留在 Stream 里的存量消息,Redis 会对比消费组last-delivered-id和 Stream 最大消息 ID,差值即为 lag;- lag 不为 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 常驻运行 + 进程托管 + 积压监控 + 幂等校验 + 死信处理。
技术进阶没有捷径,但有高效方法。 本号专注分享实战干货、避坑指南、性能调优、面试重难点, 每一篇都是亲手落地测试,帮你少走弯路、快速提升核心竞争力。
❤️ 点赞、在看、收藏、关注一键安排
关注账户,第一时间获取硬核技术干货,下期不见不散!