订单超时取消这个需求,几乎是每个电商、外卖、票务项目都会碰到的。用户下单之后 15 分钟没支付,就要把订单关掉、库存还回去、优惠券退给用户。看起来简单,做对不容易。
我前后接过四个类似的需求,用过三种方案,翻过两次车。最后一次翻得比较狠,是双十一大促当天,一批订单该取消的没取消,第二天客服电话被打爆了。事故原因说出来可能有人会笑:List 队列在高峰期拥堵,一堆过期订单排在队列尾巴上,等轮到我处理的时候早就超出「15 分钟」这个时间窗了,用户已经把钱付了,订单却被关了。赔了几千块。
那次之后我彻底把方案换成了 Redis Stream,稳定运行了大半年,分享出来给同路人。
先说清 List 队列的三个坑
PHP 生态里,用 Redis List 做队列是最常见的做法。LPUSH 塞进去,BRPOP 拿出来,加个定时任务或者常驻进程消费,好像就齐活了。但用久了会发现有三个躲不开的问题。
第一个坑:消息一旦被 BRPOP 弹出,就从 Redis 里消失了。如果消费者拿到消息还没处理完就崩了(PHP 进程 fatal error、服务器断电、被 OOM killer 干掉),这条消息就丢了,没有任何补救。List 的语义是「一次性消费」,不保证「至少一次」。
第二个坑:没有消费者组的概念。要扩容消费者,只能靠多个进程同时 BRPOP,但每个消息只能被一个消费者拿到,做负载均衡。如果某个消费者处理慢,堆到它那的消息只能干等 —— 因为 List 是队列,先进先出,头部的消息没处理完,后面的动不了。
第三个坑:延迟队列要靠两个队列倒腾。延迟消息的标准做法是:先塞到一个 ZSet,score 是到期时间戳;然后一个「搬运工」进程轮询 ZSet,把到期消息搬到 List 里,消费者再从 List 消费。这样至少有两个进程、两个数据结构在维护,但凡其中一环挂了,整条链路就断。
我上次事故就是踩在第三个坑上。搬运工进程的轮询周期设成了 5 秒,平时没问题,大促那会儿 ZSet 里堆了十几万条消息,一次搬运几百条,每分钟只能消化一部分,订单延迟取消的时间从 15 分钟漂到了 22 分钟。
为什么选 Redis Stream
Redis 5.0 开始有 Stream 这个数据结构,定位就是「消息队列」。它把前面所有痛点都解决了:
- 有消费者组概念,多个消费者可以分担同一个 Stream 的消息
- 有 ACK 机制,消息被消费但没确认之前,会存在 Pending Entries List(PEL)里,崩溃重启后能捞回来
- 支持 XAUTOCLAIM,能自动接管「其他消费者处理了很久还没 ACK」的消息,防止单点卡死
- 天然支持多消费者负载均衡
用 List 做队列,本质上是「用一个不专门做队列的数据结构硬撑队列场景」。Stream 则是官方为队列场景设计的。既然 ThinkPHP 里已经默认封装了 Redis 扩展,Stream 的 API 也能直接用,那就没道理不用。
至于延迟功能,Stream 本身不做延迟,还是要配合 ZSet。但配合的方式跟之前不一样了 —— 搬运工搬到 Stream 里的消息,会自动进入消费者组,有 ACK 保护,安全性提升一大截。
整体架构长什么样
先摆出设计图,后面代码都是照这个实现的。
| 组件 | 数据结构 | 职责 |
|---|---|---|
| 延迟池 | ZSet | 存所有定时消息,score 是到期时间戳 |
| 投递搬运工 | 常驻进程 | 轮询 ZSet,把到期消息 XADD 到 Stream |
| 消息通道 | Stream | Stream 里存待处理消息,带消费者组 |
| 消费者 | 常驻进程 | XREADGROUP 读消息,处理后 XACK 确认 |
| 兜底捞回 | 定时进程 | XAUTOCLAIM 把超时未 ACK 的消息接管回来 |
四个常驻进程的角色分工很清楚。生产环境里我会把消费者跑多份,其余三个各跑一份就够了。
核心数据结构:消息怎么编码
消息里要带什么,决定了后面消费者能做到多灵活。我的编码格式是:
{
"event": "order.timeout", // 事件类型,一个 Stream 里可以跑多种事件
"payload": { // 业务数据
"order_no": "SO20260118X0001",
"created_at": 1768763400
},
"retry": 0, // 重试次数,用于死信判断
"ts": 1768764300 // 消息投递时间,方便排查
}
event 字段很关键。一个项目里往往有十几种定时场景 —— 订单超时、售后超时、预约提醒、退款重试。如果每种都是独立 Stream,维护成本高;用 event 字段在一个 Stream 里区分,交换机一样的效果。
生产者:把消息塞进延迟池
生产者随时被业务代码调用。API 层下单成功后调一下 DelayQueue::push(),延迟消息就排上了。
<?php
namespace appcommonservice;
use thinkfacadeCache;
use thinkfacadeLog;
class DelayQueue
{
/** ZSet 的 key */
private const ZSET_KEY = 'mq:delay:pool';
/**
* 投递一条延迟消息
* @param string $event 事件类型,如 order.timeout
* @param array $payload 业务数据
* @param int $delayIn 延迟秒数
*/
public static function push(string $event, array $payload, int $delayIn): void
{
$redis = Cache::store('redis')->handler();
$message = json_encode([
'event' => $event,
'payload' => $payload,
'retry' => 0,
'ts' => time(),
], JSON_UNESCAPED_UNICODE);
$score = time() + $delayIn;
// 用唯一编号做 member,防止 member 重复导致覆盖
$member = time() . ':' . bin2hex(random_bytes(6));
$redis->zAdd(self::ZSET_KEY, $score, $member);
// 消息内容跟 member 对应关系放在 Hash 里,避免 ZSet 里塞大字符串
$redis->hSet(self::ZSET_KEY . ':data', $member, $message);
}
/**
* 取消延迟消息(比如用户下单后立刻付款了,不需要再取消了)
*/
public static function cancel(string $member): void
{
$redis = Cache::store('redis')->handler();
$redis->zRem(self::ZSET_KEY, $member);
$redis->hDel(self::ZSET_KEY . ':data', $member);
}
}
有两个设计选择说明一下。
为什么 ZSet 的 member 不用 JSON 字符串本身?因为 ZSet 的 member 越长,内存占用越大,而且 ZRANGE 返回结果要回传整个字符串给客户端,网络开销也大。用「时间戳 + 随机数」的短标识,实际数据放 Hash,两个结构的维护成本不高,性能好很多。
为什么不用 pipeline?这里每条消息一次 ZADD、一次 HSET,一共两次网络往返。业务量不大的话没问题,量大可以合并 pipeline。但要注意 pipeline 里如果有失败,回滚很麻烦,需要业务侧自己补偿。我这里保守选择了单次写,牺牲一点性能换透明性。
搬运工:把到期的消息扔进 Stream
搬运工是一个常驻进程,每秒轮询 ZSet,挑出到期的消息投递到 Stream。核心代码:
<?php
namespace appcommand;
use thinkconsoleCommand;
use thinkconsoleInput;
use thinkconsoleOutput;
use thinkfacadeCache;
use thinkfacadeLog;
class DelayMover extends Command
{
private const ZSET_KEY = 'mq:delay:pool';
private const STREAM = 'mq:events';
/** 单次最多搬运条数,防止一口气拉太多 */
private const BATCH = 200;
/** 无消息时的休眠时间(微秒) */
private const SLEEP_US = 200_000;
protected function configure()
{
$this->setName('mq:delay-mover')
->setDescription('把到期消息从 ZSet 搬到 Stream');
}
protected function execute(Input $input, Output $output)
{
$redis = Cache::store('redis')->handler();
$output->writeln('delay-mover started: ' . date('Y-m-d H:i:s'));
while (true) {
try {
$now = time();
// 拉出到期消息的 member
$members = $redis->zRangeByScore(
self::ZSET_KEY,
'-inf',
(string)$now,
['limit' => [0, self::BATCH]]
);
if (empty($members)) {
usleep(self::SLEEP_US);
continue;
}
// 一次性把消息内容捞出来
$datas = $redis->hMGet(self::ZSET_KEY . ':data', $members);
foreach ($members as $i => $member) {
$message = $datas[$i] ?? null;
if ($message === null) {
// 数据丢了,ZSet 里残留,清理掉
$redis->zRem(self::ZSET_KEY, $member);
continue;
}
// XADD 投递到 Stream,用 * 让 Redis 自动生成 ID
$redis->xAdd(self::STREAM, '*', ['msg' => $message]);
// 投递成功后清理
$redis->zRem(self::ZSET_KEY, $member);
$redis->hDel(self::ZSET_KEY . ':data', $member);
}
$count = count($members);
if ($count >= self::BATCH) {
// 一次搬运拉满了,说明堆积很多,立即处理下一批
continue;
}
usleep(self::SLEEP_US);
} catch (Throwable $e) {
Log::error('delay-mover error: ' . $e->getMessage());
usleep(1_000_000);
}
}
}
}
这段代码里有两个值得一提的点。
搬运顺序是「先 XADD,再 ZREM」。这个顺序不能反。如果先 ZREM 再 XADD,中间进程崩了,消息就彻底丢了。反过来,如果 XADD 成功但 ZREM 失败,最多是消息被重复投递一次,消费者不会重复处理(下面会讲怎么用业务 ID 去重)。宁可重复不可丢失,这是队列的基本盘。
空了就睡 200ms。这个轮询周期是可调的。太短了 CPU 空转,太长了延迟不准。200ms 是一个平衡点,订单超时这种场景完全够用,像秒杀倒计时这种高精度延迟需求就不适合了 —— 那种要用专门的定时器方案。
消费者:从消费者组里取消息
消费者进程是真正干活的地方。它做的事:从 Stream 读一条消息,按 event 分发到对应处理器,处理成功就 XACK,失败就记录并让它留在 PEL 里等待后续重试。
<?php
namespace appcommand;
use thinkconsoleCommand;
use thinkconsoleInput;
use thinkconsoleOutput;
use thinkfacadeCache;
use thinkfacadeLog;
class EventConsumer extends Command
{
private const STREAM = 'mq:events';
private const GROUP = 'main';
private const CONSUMER = 'worker'; // 稍后会加进程号
private const BLOCK_MS = 2000;
/** 消息处理超过这个秒数还没 ACK,会被其他消费者接管 */
private const ACK_TIMEOUT_MS = 60_000;
private const MAX_RETRY = 3;
protected function configure()
{
$this->setName('mq:consumer')
->setDescription('消费 Stream 里的消息');
}
protected function execute(Input $input, Output $output)
{
$redis = Cache::store('redis')->handler();
$consumer = self::CONSUMER . ':' . getmypid();
// 消费者组不存在就创建,MKSTREAM 允许 Stream 还不存在时也创建
try {
$redis->xGroup('CREATE', self::STREAM, self::GROUP, '0', true);
} catch (Throwable $e) {
// 组已存在会抛异常,忽略
}
$output->writeln("consumer started: {$consumer}");
while (true) {
try {
$result = $redis->xReadGroup(
self::GROUP,
$consumer,
[self::STREAM => '>'], // '>' 表示只读未投递给任何消费者的新消息
1, // 一次读 1 条
self::BLOCK_MS, // 阻塞等待时间
false
);
if (empty($result)) {
// 阻塞超时了,顺手处理一下 PEL 里的旧消息
$this->reclaimStale($redis, $output);
continue;
}
// result 结构: [streamName => [ id => [field => value] ]]
foreach ($result as $streamName => $messages) {
foreach ($messages as $id => $fields) {
$this->handleOne($redis, $id, $fields['msg'] ?? '');
}
}
} catch (Throwable $e) {
Log::error('consumer error: ' . $e->getMessage());
sleep(1);
}
}
}
private function handleOne($redis, string $id, string $raw): void
{
try {
$message = json_decode($raw, true);
if (!is_array($message)) {
throw new RuntimeException('message decode failed');
}
$event = $message['event'] ?? '';
$payload = $message['payload'] ?? [];
$retry = (int)($message['retry'] ?? 0);
// 按事件类型分发
$handlerClass = $this->resolveHandler($event);
if ($handlerClass === null) {
// 未知事件,直接 ACK 掉,避免堆积
$redis->xAck(self::STREAM, self::GROUP, [$id]);
Log::warning("unknown event: {$event}");
return;
}
$handler = new $handlerClass();
$handler->handle($payload);
// 处理成功,ACK
$redis->xAck(self::STREAM, self::GROUP, [$id]);
} catch (Throwable $e) {
Log::error("process failed id={$id} err=" . $e->getMessage());
$retry = (int)(json_decode($raw, true)['retry'] ?? 0) + 1;
if ($retry >= self::MAX_RETRY) {
// 超过重试上限,写死信,然后 ACK 掉防止无限重试
$this->deadLetter($redis, $raw, $e->getMessage());
$redis->xAck(self::STREAM, self::GROUP, [$id]);
}
// 未达上限就不 ACK,让这条消息留在 PEL 里,稍后被重新接管
}
}
private function resolveHandler(string $event): ?string
{
$map = [
'order.timeout' => appcommonmqhandlerOrderTimeoutHandler::class,
'aftersale.timeout' => appcommonmqhandlerAftersaleTimeoutHandler::class,
'coupon.expire' => appcommonmqhandlerCouponExpireHandler::class,
];
return $map[$event] ?? null;
}
private function deadLetter($redis, string $raw, string $reason): void
{
$redis->xAdd('mq:dead-letter', '*', [
'msg' => $raw,
'reason' => $reason,
'ts' => time(),
]);
}
/**
* 接管那些「被消费了但超时没 ACK」的消息
*/
private function reclaimStale($redis, Output $output): void
{
try {
$result = $redis->xAutoClaim(
self::STREAM,
self::GROUP,
self::CONSUMER . ':' . getmypid(),
self::ACK_TIMEOUT_MS,
'0-0',
['COUNT' => 20]
);
if (empty($result)) {
return;
}
// phpredis 的返回大致是 [nextCursor, messages] 或 [messages]
$messages = $result[1] ?? $result;
foreach ($messages as $id => $fields) {
$output->writeln("reclaim {$id}");
$this->handleOne($redis, $id, $fields['msg'] ?? '');
}
} catch (Throwable $e) {
Log::warning('xAutoClaim failed: ' . $e->getMessage());
}
}
}
几个关键设计点展开讲讲。
消费者名带 PID。多个消费者进程跑在同一台机器上时,如果消费者名都一样,Redis 会认为它们是同一个消费者,PEL 里会混在一起。加上 PID 就区分开了。
用 > 读新消息。这是 XREADGROUP 的一个特殊语法。用 > 读的是「还没投递给任何消费者的消息」,也就是新消息。如果用具体的 ID 比如 0,读的是 PEL 里还没 ACK 的消息,用于崩溃恢复。
失败不 ACK 是为了重试。处理失败的时候不调 XACK,这条消息就会留在消费者组的 PEL 里。等超过 ACK_TIMEOUT_MS(60 秒),会被 XAUTOCLAIM 抓回来重新处理。这就是「至少一次」的实现方式。
重试上限是必要的。如果某个消息因为 Bug 一直处理失败,会无限重试,把 PEL 撑爆。到 3 次就扔死信,人工介入。
业务处理器的写法:幂等是底线
既然消息可能被重复投递,业务处理器就必须幂等。这是不能省的。
<?php
namespace appcommonmqhandler;
use thinkfacadeDb;
use thinkfacadeLog;
class OrderTimeoutHandler
{
public function handle(array $payload): void
{
$orderNo = $payload['order_no'] ?? '';
if (empty($orderNo)) {
throw new InvalidArgumentException('order_no missing');
}
$order = Db::name('order')->where('order_no', $orderNo)->find();
if (!$order) {
// 订单不存在,直接返回,不需要重试
Log::info("timeout handler: order not found, skip. {$orderNo}");
return;
}
// 幂等守卫:只有未支付的订单才处理
if ($order['status'] !== 'pending') {
Log::info("timeout handler: order already {$order['status']}, skip. {$orderNo}");
return;
}
Db::startTrans();
try {
// 用条件更新保证并发安全
$affected = Db::name('order')
->where('id', $order['id'])
->where('status', 'pending')
->update([
'status' => 'cancelled',
'cancel_reason'=> 'timeout',
'cancelled_at' => time(),
'updated_at' => time(),
]);
if ($affected !== 1) {
// 状态被别人改过了,本次不处理
Db::rollback();
return;
}
// 库存归还
foreach ($order['items'] as $item) {
Db::name('sku')
->where('id', $item['sku_id'])
->inc('stock', $item['qty'])
->update();
}
// 优惠券归还
if (!empty($order['coupon_id'])) {
Db::name('coupon_user')
->where('id', $order['coupon_id'])
->update([
'status' => 'unused',
'used_at' => null,
'updated_at' => time(),
]);
}
Db::commit();
Log::info("order timeout cancelled: {$orderNo}");
} catch (Throwable $e) {
Db::rollback();
throw $e;
}
}
}
条件更新那一行是重点。判断 $affected !== 1 就提前退出,能扛住「并发消费」和「重试消费」两种重复场景。这就是之前那篇文章讲过的幂等套路,在队列消费里同样是基础操作。
进程怎么跑起来
四个常驻进程用 ThinkPHP 的命令来启动。生产环境建议用 supervisor 托管,别只靠 nohup。
# 搬运工,一个就够
php think mq:delay-mover
# 消费者,按 CPU 核数决定开几个
php think mq:consumer
supervisor 的配置大概长这样:
[program:mq-mover]
command=php /path/to/think mq:delay-mover
directory=/path/to/project
autostart=true
autorestart=true
user=www
[program:mq-consumer]
command=php /path/to/think mq:consumer
directory=/path/to/project
numprocs=4
process_name=%(program_name)s_%(process_num)02d
autostart=true
autorestart=true
user=www
注意 numprocs=4 那一行。四个消费者同时跑,会自动负载均衡,因为它们在同一个消费者组里。
观察和监控
队列跑起来之后必须能看。我一般用这几个命令,在运维机上随时敲:
# 查看 Stream 里的消息总数
redis-cli XLEN mq:events
# 查看消费者组的状态,包括 PEL 里堆积了多少
redis-cli XINFO GROUPS mq:events
# 查看具体的 PEL 内容
redis-cli XPENDING mq:events main
# 查看延迟池里还剩多少
redis-cli ZCARD mq:delay:pool
# 查看死信
redis-cli XLEN mq:dead-letter
我一般会把 XPENDING 的输出做成告警指标。如果 PEL 里的消息数超过某个阈值(比如说 100),或者最老的一条已经停留超过 10 分钟,说明消费者出问题了,得赶紧看。
压测数据
本地机器跑压力测试,环境是 4 核 8G、Redis 7.2、PHP 8.2。测试场景是每秒投递 2000 条延迟消息,延迟时间随机分布在 1 到 60 秒之间。
| 指标 | List 方案 | Stream 方案 |
|---|---|---|
| 峰值投递速率 | 3500/s | 3200/s |
| 消费者吞吐(单进程) | 1800/s | 1500/s |
| 搬运延迟中位数 | 320ms | 210ms |
| 搬运延迟 P99 | 1800ms | 480ms |
| 崩溃后消息丢失量 | 不确定(丢失) | 0 |
吞吐略低一点,因为 Stream 的每次操作会有个 XACK 的额外往返,还有 PEL 的维护开销。但延迟这一项,尤其是 P99 延迟,Stream 优势很明显。而且「崩溃后零丢失」这一条,值回所有开销。
几个还没解决的坑
说实话,这套方案也有它的边界。
第一,常驻进程对内存的占用会随时间增长。PHP 进程虽然不像 JVM 那样吃内存,但长时间运行之后,OpCache、各种框架的静态缓存,还是会慢慢涨。我的做法是监听进程内存占用,超过 200MB 就主动退出让 supervisor 拉起来。这个由 supervisor 处理比在代码里处理更省事。
第二,XDEL 清理是必要的。Stream 里的消息即使被 ACK,也不会自动删除,会一直累积。要么定期 XTRIM,要么用 MINID 参数 XADD。我一般用定时任务每天凌晨执行:XTRIM mq:events MINID ~,只保留最近一天。
第三,跨机房部署会麻烦。这套方案依赖 Redis 是单点或者主从的,如果你要跨机房,延迟池在 A 机房、消费者在 B 机房,网络抖动会导致大量重试。要么用 Redis Cluster 保证数据同区,要么把整个队列部署在同一机房内。
写在最后
Redis Stream 不是银弹,它只是把「可靠」这件事从你自己维护变成了数据结构本身提供的特性。用 List 做队列的时候,可靠性要靠业务代码去补,补得好不好全看写代码的人当时的清醒程度。凌晨三点写出来的代码和下午三点写出来的代码,能是一个水平吗。
Stream 的 ACK 和 PEL 机制,是在数据结构层面给你兜底,不管写代码的人当时是什么状态,机制都在那里。这就是我觉得它值得用的根本原因。
上面这套代码我在三个项目里用过,最长的一个跑了 14 个月没出过问题。如果你正好在做超时取消、延迟提醒、异步任务这类场景,可以直接拿去改,也可以根据自己情况调整。别的不敢说,至少比 List 稳。

