上周接了个紧急工单。用户投诉说支付成功了但优惠券没到账,客服查了后台发现订单状态是「已支付」,但券确实没发。翻日志,支付回调接口的响应时间是 3.2s,而支付平台那边的超时阈值是 3 秒。回调被判定为超时,它重试了一次,这次接口顺利返回,但第一次其实也执行完了——于是券发了两张。
问题挺典型的。一个接口做了太多事:改订单状态、发优惠券、发短信通知、更新用户积分、往统计表里写一条记录。数据库操作串起来大概 3 秒,赶上支付平台的打款延迟,偶尔就超了。
处理方案是标准的:把「必须同步完成」和「可以稍后完成」的两类操作拆开。订单状态必须同步改,其余的都可以异步。
先看现状:一个回调查了多少张表
原来的回调处理大概是这个结构:
public function notify()
{
$params = $this->request->post();
$sign = $this->verifySign($params);
if (!$sign) {
return json(['code' => 'FAIL']);
}
Db::startTrans();
try {
// 1. 更新订单状态
Order::where('order_no', $params['out_trade_no'])
->update(['status' => 'paid', 'paid_at' => time()]);
// 2. 发优惠券
$this->grantCoupon($params['out_trade_no']);
// 3. 更新用户积分
$this->updatePoints($params['out_trade_no']);
// 4. 发短信通知
$this->sendSms($params['out_trade_no']);
// 5. 写统计
$this->writeStatistics($params['out_trade_no']);
Db::commit();
} catch (Throwable $e) {
Db::rollback();
return json(['code' => 'FAIL']);
}
return json(['code' => 'SUCCESS']);
}
这个写法有两个层次的问题。
一个是性能问题。五件事串行执行,其中短信发送是走第三方接口的,光是这一项就可能花掉 500 毫秒到 1 秒。整体 3 秒是合理的。
另一个是设计问题。「更新订单状态」和「发短信」的重要级别完全不一样。前者失败业务就断了,必须管;后者失败了顶多用户晚一分钟收到通知,不该让整个回调失败甚至回滚订单状态。全部塞在一个事务里,等于让一个小概率的第三方超时拖垮整个支付流程。
为什么选了 think-queue
异步化的方案不止一种。简单说一下当时的取舍:
- 扔给系统 crontab 定时扫表。最简单,但延迟最差。优惠券要等下一分钟才到,用户可能会投诉。而且数据量大之后,扫表本身就是个问题。
- 上 RabbitMQ 或 Kafka。功能全,但对这个小项目来说是杀鸡用牛刀。团队里没人维护过,一旦出问题响应会慢。
- Swoole 常驻内存。整个项目得改造成常驻模式,这个改动面太大了,不是一次故障就能推动的。
- think-queue + Redis。装一个包、配一下驱动就能跑,学习成本几乎为零。Redis 本来就是项目里跑着的组件,不用新引入依赖。
选了第四个。它不完美——比如没有 RabbitMQ 那么强的消息可靠性保证,但对于「发短信、发券、写统计」这类允许少量丢失、允许重试的场景,够了。
安装和配置
装包:
composer require topthink/think-queue
配置文件在 config/queue.php。默认会给一份示例,改一下连接信息:
<?php
return [
'default' => 'redis',
'connections' => [
'redis' => [
'type' => 'redis',
'queue' => 'default',
'host' => '127.0.0.1',
'port' => 6379,
'password' => '',
'select' => 0,
'timeout' => 0,
'persistent' => false,
],
],
'failed' => [
'type' => 'none',
'table' => 'failed_jobs',
],
];
把 failed 的 type 改成 database,失败的任务就会落库,方便排查。对应的表通过命令生成:
php think queue:table
php think migrate:run
把发短信这个动作拆出来
先做最简单的那个。新建一个 Job 类,放在 app/job/SendSmsJob.php:
<?php
namespace appjob;
use thinkqueueJob;
use thinkfacadeLog;
class SendSmsJob
{
public function fire(Job $job, array $data): void
{
$orderNo = $data['order_no'] ?? '';
try {
$ok = $this->doSend($orderNo);
} catch (Throwable $e) {
Log::error('短信发送异常', ['order_no' => $orderNo, 'msg' => $e->getMessage()]);
$ok = false;
}
if ($ok) {
$job->delete();
return;
}
if ($job->attempts() > 3) {
Log::error('短信任务失败次数过多,转入失败队列', ['order_no' => $orderNo]);
$job->delete();
return;
}
// 10 秒后重试
$job->release(10);
}
private function doSend(string $orderNo): bool
{
$order = appmodelOrder::where('order_no', $orderNo)->find();
if (!$order) {
return false;
}
// 调用第三方短信接口
return app(appserviceSmsService::class)
->sendPaidNotice($order->user_id, $orderNo);
}
}
几个地方值得停一下。
fire() 是 think-queue 约定的入口。第二个参数就是推任务时传的数据,比在任务里重新查一遍要高效,但注意这里只放必要标识(order_no),别把一整条订单对象序列化进去——订单字段后续变了,老任务会拿到过期数据。
$job->attempts() 是当前已经尝试了几次。写在 > 3 是因为 think-queue 的 attempts 从 1 开始。超过上限时不要直接释放,否则任务会一直转,要有明确的退出条件。
release(10) 表示延迟 10 秒重投。不是立即重试,是因为第三方接口这类故障通常是瞬时的(限流、抖动),立即重试大概率还会失败,反而把队列堵住。10 秒是个经验值,具体看业务。
发券要保证事务和幂等
发券比发短信复杂一点——它涉及数据库写操作,而且绝对不能重复发。
<?php
namespace appjob;
use thinkqueueJob;
use thinkfacadeDb;
use thinkfacadeLog;
class GrantCouponJob
{
public function fire(Job $job, array $data): void
{
$orderNo = $data['order_no'] ?? '';
try {
Db::startTrans();
// 用唯一索引兜底,防止并发或重试导致重复发券
$exists = Db::table('coupon_log')
->where('order_no', $orderNo)
->lock(true)
->find();
if ($exists) {
Db::commit();
$job->delete();
return;
}
$order = Db::table('order')->where('order_no', $orderNo)->find();
if (!$order) {
Db::rollback();
$job->delete();
return;
}
$coupon = $this->pickCoupon($order);
if ($coupon) {
Db::table('user_coupon')->insert([
'user_id' => $order['user_id'],
'coupon_id' => $coupon['id'],
'order_no' => $orderNo,
'expire_at' => time() + $coupon['valid_days'] * 86400,
'created_at' => time(),
]);
Db::table('coupon_log')->insert([
'order_no' => $orderNo,
'coupon_id' => $coupon['id'],
'created_at' => time(),
]);
}
Db::commit();
$job->delete();
} catch (Throwable $e) {
Db::rollback();
Log::error('发券任务异常', ['order_no' => $orderNo, 'msg' => $e->getMessage()]);
if ($job->attempts() > 3) {
$job->delete();
return;
}
$job->release(30);
}
}
private function pickCoupon(array $order): ?array
{
// 根据订单金额选一张合适的券
return Db::table('coupon')
->where('min_amount', '<=', $order['amount'])
->where('status', 1)
->order('min_amount', 'desc')
->find() ?: null;
}
}
lock(true) 是悲观锁(对应 SELECT ... FOR UPDATE),加在 coupon_log 的查询上。可能你会想为什么不直接在 user_coupon 上用唯一索引,那样更简单。真实考虑是:同一张券可以发给不同用户,唯一键不好设计;但「同一个订单只能发一次券」这个约束是明确的,落在 coupon_log 上最合适。
另外 coupon_log 表在生产环境必须加唯一索引:
ALTER TABLE `coupon_log`
ADD UNIQUE INDEX `uk_order_no` (`order_no`);
代码里的 lock(true) 是应用层的锁,唯一索引是数据库层的兜底。两者都要有,前者防并发,后者防意外(比如有人手动执行了 job 代码)。
改造回调接口
现在原来的 notify() 可以简化了。同步的只剩订单状态和拆解后的任务推送:
public function notify()
{
$params = $this->request->post();
if (!$this->verifySign($params)) {
return json(['code' => 'FAIL']);
}
$orderNo = $params['out_trade_no'];
// 只在订单状态未变时才更新,天生幂等
$affected = Order::where('order_no', $orderNo)
->where('status', '<>', 'paid')
->update(['status' => 'paid', 'paid_at' => time()]);
if ($affected === 0) {
// 说明是重复回调,直接返回成功
return json(['code' => 'SUCCESS']);
}
// 拆成三个独立任务
Queue::push(GrantCouponJob::class, ['order_no' => $orderNo], 'order');
Queue::push(SendSmsJob::class, ['order_no' => $orderNo], 'notify');
Queue::push(StatisticsJob::class, ['order_no' => $orderNo], 'stats');
return json(['code' => 'SUCCESS']);
}
这里有两个关键设计。
「只更新未支付的订单」这个条件天然防重。支付平台重复回调时,第一次把状态改为 paid,第二次的 update 就影响 0 行,直接返回成功。不用额外做去重表,数据库自己的行锁就能保证。
三个任务放到三个不同的队列名里。因为它们的处理速度差很多:发券大概是毫秒级,写统计可能更慢,但发短信是走第三方的,动不动一两秒。如果都放默认队列,一个慢短信会堵住后面所有的发券。分开之后可以起多个消费进程,各跑各的。
起消费进程
开发环境用 queue:listen:
php think queue:listen --queue order,notify,stats
listen 每次处理完一个任务会重启一次进程,方便改代码立即生效。生产环境必须用 work:
php think queue:work --queue order --daemon \
--tries=3 --timeout=60 --memory=256
四个参数逐个说:
--daemon:不重启进程,处理完继续等下一条。性能是 listen 的好几倍。--tries=3:最多尝试三次,超过之后(如果配了失败队列)落库。--timeout=60:单条任务最长执行 60 秒,超了会被强制中断并记一次尝试。--memory=256:进程内存超过 256MB 自动退出重起,防止内存泄漏。这个参数是生产环境的救命稻草,一定要加上。
多进程部署的话,用 supervisor 起 4 个 order 队列的进程、2 个 notify 的进程、1 个 stats 的进程,根据任务量调整。
延迟任务:15 分钟后催评价
think-queue 支持延迟推任务。用户支付 15 分钟后推送一条「请给订单评价」的短信:
Queue::later(900, SendReviewSmsJob::class, ['order_no' => $orderNo], 'notify');
900 是秒数,也就是 15 分钟。
这里有个认知上的坑:延迟队列并不是「15 分钟后再执行」这么精确。Redis 的延迟实现是把任务放在一个 zset 里,score 是触发时间戳,写一个消费进程定期扫描到期的任务再挪到真正的 job 队列里。所以实际延迟 = 900 秒 + 扫描周期 + 队列排队时间。
另外还有一个更隐蔽的问题:延迟任务推出去之后,用户可能会取消订单。任务到了执行时间点,发现订单已经不在了。所以延迟任务的入口处必须重新查一遍业务状态。
public function fire(Job $job, array $data): void
{
$order = Order::where('order_no', $data['order_no'])->find();
// 延迟任务的常见坑:到点执行时,业务状态可能已经变了
if (!$order || $order['status'] !== 'paid') {
$job->delete();
return;
}
// 正常发送
}
失败任务的处理
配置了 failed 落库之后,可以在后台做一个小页面,把失败的任务列表展示出来,支持一键重投:
php think queue:failed # 查看所有失败任务
php think queue:retry {id} # 重投指定任务
php think queue:forget {id} # 删除指定任务
php think queue:flush # 清空所有失败任务
这四条命令日常排查时会经常用到。特别是 queue:failed,看一眼失败任务的数量,大概就知道系统是不是出问题了。
还有一个容易被忽略的点:失败任务积累地很快。如果不定期清理,几个月后 failed_jobs 表可能几十万条。我的做法是写一个每小时跑的清理任务,删除 30 天之前的记录,同时给团队发一条告警——因为失败任务的存在本身就是一个需要关注的信号。
监控和告警
队列天生适合暗地里坏——用户看到「支付成功」,但优惠券一直没到,就要等客服找上门才知道。所以监控是必须的。
最朴素的做法是每分钟检查一次指标,异常时报警:
# 用 redis-cli 查看队列长度
redis-cli llen queues:order
redis-cli llen queues:notify
redis-cli zcard queues:order:delayed
如果 queues:order 的长度持续几分钟超过 1000,说明消费速度已经跟不上生产速度了。这时候该做的就是加消费进程,或者查一下是不是任务执行变慢了。
另一个更有价值的信号是「成功支付但迟迟没有券」这类业务指标的对比。用一个定时任务扫最近 5 分钟的已支付订单,看看它们的发券日志有没有全部生成。这能抓出来一些队列系统本身发现不了的问题,比如任务被无辜丢弃。
五个踩过的坑
一、任务里没 catch 异常,队列直接崩。think-queue 的 fire() 方法里如果抛了没被捕获的异常,整个消费进程会退出。用 supervisor 可以自动拉起来,但你会发现在那几秒里所有任务都卡住了。业务代码里每一段独立的操作都要包 try/catch,然后把是否需要重投交给 $job->release() 决定。
二、Job 类里依赖注入不生效。think-queue 处理任务时是把 Job 类 new 出来的,不走容器的依赖注入解析。如果你想在构造函数里注入 Service,会失败。正确的做法是在 fire() 方法里用 app() 或 Container::get() 手动取。
public function fire(Job $job, array $data)
{
$service = app(appserviceCouponService::class);
// ...
}
三、任务数据里的时间字段不要用相对值。比如任务是「10 分钟后发短信」,别把 'delay' => 600 塞进数据里让 Job 自己去 sleep。任务数据只放业务标识,延迟交给 Queue::later 处理。原因是一旦任务在队列里堆积,sleep 的时间会更长,和实际业务语义完全脱节。
四、不要在任务里写「只执行一次」的代码。队列系统本质上不保证「只执行一次」——它保证的是「至少执行一次」。同一个任务可能因为超时被重投两次,也可能因为进程重启被恢复。所以任务逻辑必须自己保证幂等,唯一索引、状态判断、乐观锁都是常见手段。
五、--tries 的计数和业务重试是两回事。框架的 --tries 会在尝试次数用尽时把任务丢进 failed 表,但如果你在业务代码里做了「失败后重新入队」,那就绕过了这个限制。老实说,进入无限重试队列的任务基本上都是 bug,不要靠它来兜底。
什么时候不该上队列
队列是个好工具,但也不是所有地方都该用。
「用户注册后发一条欢迎邮件」这种操作,如果同步执行只要 200 毫秒,也没必要非得上异步——一次请求里多 200 毫秒和少 200 毫秒,用户感受不出来。引入队列反而多了一层间接和排查成本。
「下单时扣库存」这种必须严格顺序、不能失败的操作,属于核心链路,不能进队列。队列的价值在于「可以晚一点、可以重试」,如果业务本身不允许晚、不允许歧义,就得同步做。
「用户申请退款后 3 天自动通过」这种跨越几天的延迟,也不是队列的长项。延迟队列的精度和存储成本都不适合这种量级,用定时任务扫表或者专门的时间轮组件更合适。
判断标准其实就一条:如果这件事慢了半拍,用户会不会感觉不对?不会的,就进队列;会的,就同步做。就这么简单。
最后看效果
改造完成后,支付回调接口的耗时从 3.2s 降到了 180ms 左右,其中大部分是签名校验和数据库更新消耗的。发短信、发券、写统计都进了队列,平均 1 秒内处理完,最慢的短信任务也稳定在 3 秒以内。
支付平台那边不再报超时,重复回调的问题自然也就消失了——第二次回调走到 $affected === 0 的时候就返回了,根本不会触发后续逻辑。
现在回头看,这次故障其实值了。如果没出事,那个把所有操作塞在一个事务里的回调可能还会挂在那里跑一两年,直到某天第三方接口再慢一点,又崩一次。故障是个提醒,不是骂人。

