别再拿 Redis List 当队列了:ThinkPHP 8 下单超时取消的 Stream 完整落地

2026-10-04 0 677

订单超时取消这个需求,几乎是每个电商、外卖、票务项目都会碰到的。用户下单之后 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 稳。

别再拿 Redis List 当队列了:ThinkPHP 8 下单超时取消的 Stream 完整落地
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 thinkphp 别再拿 Redis List 当队列了:ThinkPHP 8 下单超时取消的 Stream 完整落地 https://www.taomawang.com/server/thinkphp/2889.html

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务