先说背景。一个跑了大概半年的商城项目,订单超时关闭最早是这么做的: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 延迟任务,代码量其实没少多少,但数据库的压力是实打实降下来了,订单关闭的时效性也从“最多延迟一分钟”变成了“基本准时”。真正需要花心思设计的是幂等——队列这个组件本身很简单,难的是让业务逻辑在“可能执行两次”的前提下依然正确。

