ThinkPHP 8 队列实战:Redis 延迟任务重做订单超时关闭与重复消费处理

2026-09-18 0 625

先说背景。一个跑了大概半年的商城项目,订单超时关闭最早是这么做的:crontab 每分钟跑一次脚本,扫一遍订单表。

SELECT * FROM tp_order
WHERE status = 0 AND create_time < UNIX_TIMESTAMP() - 1800
LIMIT 500;

订单表到 300 万行的时候问题就出来了。create_time 这个索引选择性太差,每分钟一次扫描,数据库 CPU 曲线跟心电图一样。高峰期一次扫 500 条还不够用,低峰期又在空转。更麻烦的是,关闭订单之后还要回滚库存、退优惠券、通知客服,脚本从最初 30 行写到了 200 多行,逻辑全挤在一起。

后来把这个逻辑迁到了 think-queue 的延迟任务上:下单的时候投一条 30 分钟后执行的任务,时间到了 worker 自己去处理,中间不需要数据库轮询参与。整个过程踩了一些坑,记录下来。

一、先把 think-queue 的运行模型弄清楚

think-queue 里有三个角色:

  • 投递方,也就是跑在 PHP-FPM 里的业务代码
  • 存储层,可以是 Redis,也可以是 MySQL 表
  • 消费方,也就是 php think queue:work 启动的 CLI 常驻进程

一次完整的链路是这样的:业务代码调用 Queue::later(1800, CloseOrder::class, ['order_id' => 123], 'order'),任务体被 JSON 编码后写进 Redis;常驻的 worker 进程从队列里取任务,用容器把 CloseOrder 实例化,然后调用它的 fire() 方法。

延迟任务在 Redis 驱动下并不是 sleep 等待,而是被塞进一个 zset。key 的结构大致是这样:

  • queues:order,主队列,list 结构,等着被消费的任务
  • queues:order:delayed,延迟队列,zset 结构,score 是任务的执行时间戳
  • queues:order:reserved,执行中的任务,zset 结构

worker 每次取任务之前,会先把 delayed 里 score 小于当前时间的任务搬到主队列。这句话请记一下,后面排查“延迟任务为什么不执行”的时候全靠它。

有人会问,数据库驱动也能用,为什么非要上 Redis。简单对比一下:数据库驱动不需要额外服务,靠 reserved_at 字段加事务行锁来抢占任务,量小的时候够用,但每次取任务都要查库,队列一长就会拖累主库;Redis 驱动的读写都在内存里,延迟任务原生支持,代价是多了一个需要运维的组件。订单量超过每天几万单,基本就该选 Redis 了。

二、安装与配置

composer require topthink/think-queue

ThinkPHP 8 对应的是 think-queue 3.x,要求 PHP 8.0 以上。装完之后在 config/queue.php 里配置连接:

<?php

return [
    'default'     => 'redis',
    'connections' => [
        'sync'  => [
            'type' => 'sync',
        ],
        'redis' => [
            'type'       => 'redis',
            'queue'      => 'default',
            'host'       => env('redis.host', '127.0.0.1'),
            'port'       => env('redis.port', 6379),
            'password'   => env('redis.password', ''),
            'select'     => 0,
            'timeout'    => 0,
            'persistent' => false,
        ],
    ],
    'failed'      => [
        'type' => 'none',
    ],
];

这里 failed 我特意配成了 none。think-queue 自带的失败任务记录比较“粗”,只存原始 payload 和异常堆栈,业务上想按订单号搜、想重新投递,都不方便。所以失败记录我自己建表控制:

CREATE TABLE `tp_job_failed_log` (
  `id`          bigint unsigned NOT NULL AUTO_INCREMENT,
  `job`         varchar(120)    NOT NULL DEFAULT '',
  `data`        varchar(1000)   NOT NULL DEFAULT '',
  `message`     varchar(2000)   NOT NULL DEFAULT '',
  `create_time` int unsigned    NOT NULL DEFAULT 0,
  PRIMARY KEY (`id`),
  KEY `idx_create_time` (`create_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

三、案例一:订单超时自动关闭

假设订单表结构是这样的:

CREATE TABLE `tp_order` (
  `id`          bigint unsigned NOT NULL AUTO_INCREMENT,
  `order_no`    varchar(32)     NOT NULL,
  `user_id`     bigint unsigned NOT NULL,
  `amount`      decimal(10,2)   NOT NULL DEFAULT '0.00',
  `status`      tinyint         NOT NULL DEFAULT 0 COMMENT '0待支付 1已支付 2已关闭',
  `create_time` int unsigned    NOT NULL DEFAULT 0,
  `close_time`  int unsigned    NOT NULL DEFAULT 0,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_order_no` (`order_no`),
  KEY `idx_status_create` (`status`, `create_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

下单接口里,订单写库成功之后再投任务。投递动作一定要放在事务外面,否则事务回滚了 Redis 里还留着一条脏任务,三十秒后 worker 去关一个不存在的订单:

<?php
declare(strict_types=1);

namespace appcommonservice;

use appjobCloseOrder;
use appmodelOrder;
use thinkfacadeQueue;

class OrderService
{
    public function createOrder(array $params): array
    {
        $order = Order::create([
            'order_no'    => $this->buildOrderNo(),
            'user_id'     => $params['user_id'],
            'amount'      => $params['amount'],
            'status'      => 0,
            'create_time' => time(),
        ]);

        // 30 分钟 = 1800 秒,投到名为 order 的队列
        Queue::later(1800, CloseOrder::class, [
            'order_id' => $order->id,
            'order_no' => $order->order_no,
        ], 'order');

        return $order->toArray();
    }
}

然后是 Job 类本体,这是整篇文章最核心的一段代码:

<?php
declare(strict_types=1);

namespace appjob;

use appcommonserviceStockService;
use thinkfacadeCache;
use thinkfacadeDb;
use thinkfacadeLog;
use thinkqueueJob;

class CloseOrder
{
    public function fire(Job $job, $data)
    {
        $orderId = (int) ($data['order_id'] ?? 0);
        if ($orderId <= 0) {
            $job->delete();
            return;
        }

        $redis   = Cache::store('redis')->handler();
        $lockKey = 'lock:close_order:' . $orderId;

        // 一把 30 秒的锁,减少同一订单被两个 worker 同时处理的概率
        if (!$redis->set($lockKey, 1, ['nx', 'ex' => 30])) {
            $job->release(5);
            return;
        }

        try {
            $affected = Db::name('order')
                ->where('id', $orderId)
                ->where('status', 0)                        // 只有待支付才能关
                ->where('create_time', '<=', time() - 1800)  // 没到时间就不动
                ->update([
                    'status'     => 2,
                    'close_time' => time(),
                ]);

            if ($affected > 0) {
                app(StockService::class)->releaseByOrder($orderId);
                Log::info('[CloseOrder] 已关闭订单 ' . $orderId);
            }

            $job->delete();
        } catch (Throwable $e) {
            Log::error('[CloseOrder] 执行异常', [
                'order_id' => $orderId,
                'attempts' => $job->attempts(),
                'message'  => $e->getMessage(),
            ]);

            if ($job->attempts() >= 3) {
                $this->recordFailure($data, $e);
                $job->delete();
                return;
            }

            $job->release(20);
        } finally {
            $redis->del($lockKey);
        }
    }

    protected function recordFailure(array $data, Throwable $e): void
    {
        Db::name('job_failed_log')->insert([
            'job'         => self::class,
            'data'        => json_encode($data, JSON_UNESCAPED_UNICODE),
            'message'     => mb_substr($e->getMessage(), 0, 500),
            'create_time' => time(),
        ]);
    }
}

这段代码里有几个点值得单独拿出来说。

1. 更新语句里的 where 条件就是幂等的锚点

where('status', 0) 不只是业务判断,它同时承担了幂等的职责。任务如果被重复执行,第二次 update 影响行数会是 0,后面的回滚库存自然也不会跑。用户自己取消了订单、或者已经支付成功了,这条任务就变成一个空操作。

Redis 那把锁的作用只是削峰,减少并发冲突。别指望它保证正确性——锁过期、finally 里提前释放,都可能让两个进程同时进来。真正兜底的是数据库条件更新。

2. create_time 的时间兜底不是多余的

我第一次上线的时候,测试环境把延迟参数写成了 30,以为是分钟,其实是秒。结果一晚上关掉两千多个正常待支付订单,客服电话被打爆。后来加了 where('create_time', '<=', time() - 1800) 这条保险:任务提前跑、人为重放、历史脏数据,都会在这里被拦住。

3. 支付和关闭之间的竞态

用户在 29 分 58 秒点了支付,关闭任务在 30 分整执行,两个请求几乎同时到达。处理办法是在支付回调里也带上状态条件:

$affected = Db::name('order')
    ->where('id', $orderId)
    ->where('status', 0)
    ->update(['status' => 1, 'pay_time' => time()]);

if ($affected === 0) {
    // 订单已经被关闭了,走退款流程
}

谁先抢到 status 从 0 变成非 0,谁就算赢,另一方拿到影响行数 0 自行处理。这就是最朴素的乐观锁,比加悲观锁简单得多。

4. 重试逻辑自己控制,别和框架混用

think-queue 的 Job 基类本身也有重试机制:任务抛出异常后,worker 会根据 --tries 参数决定重新入队还是丢弃。但我在 fire 里把所有异常都 catch 了,框架那层就拿不到异常,也就不会介入。两条路只能选一条,混着用会出现“重试了 3 次又重试 3 次”的尴尬局面。

四、案例二:支付成功后的异步派发

用户支付成功后要做的事情很多:加积分、发优惠券、通知商家、更新商品销量。这些操作如果全写在支付回调接口里,接口响应时间轻松超过三秒,第三方支付平台会一直重推通知。所以回调里只做状态变更,其余全部投任务。

<?php

use appjobGrantBenefit;
use thinkfacadeDb;
use thinkfacadeQueue;

public function onPaid(string $orderNo): void
{
    $order = Db::name('order')->where('order_no', $orderNo)->find();
    if (!$order) {
        return;
    }

    // 先改状态,用影响行数保证回调幂等
    $affected = Db::name('order')
        ->where('id', $order['id'])
        ->where('status', 0)
        ->update(['status' => 1, 'pay_time' => time()]);

    if ($affected === 0) {
        // 支付平台重复推送,直接返回成功
        return;
    }

    Queue::push(GrantBenefit::class, ['order_id' => $order['id']], 'benefit');
}

派发权益的任务里,用一张日志表加唯一索引来做幂等:

CREATE TABLE `tp_order_benefit_log` (
  `id`           bigint unsigned NOT NULL AUTO_INCREMENT,
  `order_id`     bigint unsigned NOT NULL,
  `benefit_type` varchar(32)     NOT NULL,
  `create_time`  int unsigned    NOT NULL DEFAULT 0,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_order_benefit` (`order_id`, `benefit_type`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
<?php
declare(strict_types=1);

namespace appjob;

use thinkdbexceptionPDOException;
use thinkfacadeDb;
use thinkqueueJob;

class GrantBenefit
{
    public function fire(Job $job, $data)
    {
        $orderId = (int) ($data['order_id'] ?? 0);
        $order   = Db::name('order')->where('id', $orderId)->find();

        // 任务自身也要校验业务状态,支付没成功一律不发
        if (!$order || (int) $order['status'] !== 1) {
            $job->delete();
            return;
        }

        try {
            foreach (['point', 'coupon', 'notify'] as $type) {
                $this->grant($order, $type);
            }
            $job->delete();
        } catch (Throwable $e) {
            if ($job->attempts() >= 3) {
                Db::name('job_failed_log')->insert([
                    'job'         => self::class,
                    'data'        => json_encode($data, JSON_UNESCAPED_UNICODE),
                    'message'     => mb_substr($e->getMessage(), 0, 500),
                    'create_time' => time(),
                ]);
                $job->delete();
                return;
            }
            $job->release(10);
        }
    }

    protected function grant(array $order, string $type): void
    {
        try {
            // 先占坑,插入成功才执行真正的发放逻辑
            Db::name('order_benefit_log')->insert([
                'order_id'     => $order['id'],
                'benefit_type' => $type,
                'create_time'  => time(),
            ]);
        } catch (PDOException $e) {
            // 1062 是唯一键冲突,说明这个权益已经发过了
            if (str_contains($e->getMessage(), '1062')) {
                return;
            }
            throw $e;
        }

        // 真正发放积分 / 优惠券的逻辑
        // ...
    }
}

这种“先写日志再执行”的顺序比“先查后写”可靠。两个 worker 同时执行,唯一索引只会让其中一个 insert 成功,另一个直接返回,不需要任何分布式锁。

五、重复消费到底从哪里来

把来源拆开看,一共三条路。

第一条,业务重试。任务里自己调了 release,或者抛异常后被框架重新入队。这条最明显,也最容易处理。

第二条,重复投递。用户连点两次下单按钮,或者支付平台推了三次通知,同一份业务数据进了两次队列。用唯一索引或状态条件更新就能挡住。

第三条,worker 被 kill。任务从主队列 pop 出来之后,会先被放进 queues:order:reserved。如果 worker 正好在执行过程中被 OOM Killer 干掉,这条任务不会消失,它会一直躺在 reserved 里。等下一个存活的 worker 来取任务,触发一次迁移,把它重新搬回主队列再执行一遍。

第三条最容易被忽略,因为它在本地开发环境几乎不会发生。所以写 Job 的时候脑子里要有一个前提:同一份数据至少会被执行两次。所有的幂等设计都要落在持久层上,Redis 锁只能减少并发,解决不了重复执行。

六、worker 怎么跑才稳

本地调试用这两个命令:

php think queue:work --queue order
php think queue:listen --queue order

listen 每次取到任务都会重新走一遍应用初始化流程,改完 Job 代码不用手动重启,适合开发阶段。生产环境请用 work –daemon。

[program:tp-order-queue]
command=php /www/shop/think queue:work --queue order --daemon --tries=3 --delay=10
directory=/www/shop
user=www
numprocs=4
autostart=true
autorestart=true
startsecs=3
stopwaitsecs=30
redirect_stderr=true
stdout_logfile=/var/log/supervisor/tp-order-queue.log
stdout_logfile_maxbytes=50MB
stdout_logfile_backups=5

几个容易踩的点:

  • --daemon 必须加。不加的话进程处理完一条任务就退出了,supervisor 会不停重启,日志里刷一堆启动记录。
  • numprocs 不是越多越好。订单关闭这类任务内部有数据库写操作,4 个进程已经能撑住每天十万单的量级。盲目开到 20 个,反而会把 MySQL 连接池占满。
  • 改完 Job 代码必须重启 worker。daemon 模式下代码是常驻内存的,不重启就永远跑旧逻辑。这个坑线上经常出现:开发改完代码推上去,观察半天发现行为没变,其实是旧进程还在。
  • 建议每天凌晨用 cron 重启一次 worker 进程。PHP 常驻进程难免有内存缓慢增长,定时重启是最省事的兜底手段。

七、队列积压了怎么查

Redis 里的 key 结构很直观,直接看就行:

# 主队列积压数量
redis-cli llen queues:order

# 延迟任务数量
redis-cli zcard queues:order:delayed

# 执行中的任务数量
redis-cli zcard queues:order:reserved

现象一:llen 一直涨,日志里却没有 worker 的输出。

先确认 worker 进程还活着,再确认投递时用的队列名和消费时指定的队列名是不是同一个。有人写 Queue::push($job, $data, 'benefit'),投到了 queues:benefit,而 worker 只跑 --queue order,任务永远不会被消费。这是“投递成功但没人处理”最常见的原因。

现象二:zcard delayed 的数量一直不降。

回看第一节那句话:delayed 里的任务只有在 worker 取任务时才会被搬到主队列。所有 worker 都挂了的情况下,你投再多延迟任务也没用,它们会一直安静地待在 zset 里。延迟任务占比高的系统,必须保证至少一个 worker 常驻。

现象三:zcard reserved 一直不降。

说明有 worker 在处理任务时挂掉了,任务还占在 reserved 里。把 worker 拉起来之后它们会自动搬回去重新执行。但如果某条任务本身就会让进程崩溃(比如处理超大附件导致 OOM),它会陷入崩溃—重启—再崩溃的死循环,需要手工清理:

redis-cli zrange queues:order:reserved 0 5

把有问题的 payload 找出来,用 zrem 删掉,再单独处理。

八、写个监控命令,别等出事了才发现

与其每次手动敲 redis-cli,不如写个自定义命令挂在 crontab 里:

<?php
declare(strict_types=1);

namespace appcommand;

use thinkconsoleCommand;
use thinkconsoleInput;
use thinkconsoleOutput;
use thinkfacadeCache;

class QueueMonitor extends Command
{
    protected function configure()
    {
        $this->setName('queue:monitor')
            ->setDescription('输出各队列积压情况,超过阈值触发告警');
    }

    protected function execute(Input $input, Output $output)
    {
        $queues = ['order', 'benefit', 'settle'];
        $redis  = Cache::store('redis')->handler();
        $alert  = [];

        foreach ($queues as $name) {
            $key      = 'queues:' . $name;
            $ready    = (int) $redis->llen($key);
            $delayed  = (int) $redis->zcard($key . ':delayed');
            $reserved = (int) $redis->zcard($key . ':reserved');

            $output->writeln(sprintf(
                '%s ready=%d delayed=%d reserved=%d',
                str_pad($name, 10),
                $ready,
                $delayed,
                $reserved
            ));

            if ($ready > 1000 || $reserved > 200) {
                $alert[] = sprintf('%s ready=%d reserved=%d', $name, $ready, $reserved);
            }
        }

        if ($alert) {
            $this->notify(implode("n", $alert));
        }

        return 0;
    }

    protected function notify(string $text): void
    {
        // 这里换成你自己的告警通道,钉钉、飞书、企业微信都可以
        // 用 Guzzle 发一条 webhook 消息即可
    }
}

然后在 config/console.php 里注册:

return [
    'commands' => [
        'queue:monitor' => appcommandQueueMonitor::class,
    ],
];

挂到 crontab 里每分钟跑一次:

* * * * * cd /www/shop && php think queue:monitor >> /var/log/tp-queue-monitor.log 2>&1

九、几条踩过坑之后的经验

  • Job 的 $data 里只存主键,不要存整个模型数组。payload 里的数据会过期,执行的时候应该以数据库里的最新状态为准。
  • 不要在 Job 里使用 request()session()cookie()。CLI 环境下没有这些上下文,直接抛异常。
  • 一个 Job 只做一件事。把加积分、发券、发短信全塞在一条任务里,出问题的时候你根本不知道重试会重复执行哪一步。
  • 不要在事务里投递任务,也不要把投递放在可能被回滚的位置。任务和数据库事务之间没有原子性,只能靠幂等去兜。
  • 生产环境别用 sync 驱动。那是同步执行,队列的意义就没了,本地测试图省事可以用用。
  • 长任务要拆。一条任务跑五分钟,worker 的并发能力就废了,而且一旦中途失败,重试的代价非常高。

回过头看,从数据库轮询换到 Redis 延迟任务,代码量其实没少多少,但数据库的压力是实打实降下来了,订单关闭的时效性也从“最多延迟一分钟”变成了“基本准时”。真正需要花心思设计的是幂等——队列这个组件本身很简单,难的是让业务逻辑在“可能执行两次”的前提下依然正确。

ThinkPHP 8 队列实战:Redis 延迟任务重做订单超时关闭与重复消费处理
收藏 (0) 打赏

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

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

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

淘吗网 thinkphp ThinkPHP 8 队列实战:Redis 延迟任务重做订单超时关闭与重复消费处理 https://www.taomawang.com/server/thinkphp/2772.html

下一篇:

已经没有下一篇了!

常见问题

相关文章

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

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