先说一个之前维护的订单系统。
下单接口最早只有一件事:写订单、扣库存。后来需求一条一条往上加——下单成功要给用户发短信;要加积分;要通知仓库;要在统计表里累计一份数据;要做一次风控快照。半年下来,控制器里的 create() 方法变成了两百多行,中间还穿插了 try-catch。
这种 Controller 长得像流水账的代码,问题不在长,而在责任混乱。订单本身”下单成功”这件事跟”发短信”没有业务上的强关系,但代码里被硬绑在一起。短信发了,用户下单失败;库存扣了,订单写进去的时候主键冲突;统计服务短时间不可用,整个下单接口全部 500——都是这种情况的真实后果。
这篇文章讲怎么用 ThinkPHP 8 的事件系统把这个模块重新组织一遍。中间会经过三次演进:先把副作用抽成事件监听器,再把监听器拆成同步的、异步的两类,最后处理队列的可靠性和幂等。每一步的动机和副作用都会讲清楚。
一、原始代码长什么样
先看看问题现场。大概是这样:
class OrderController
{
public function create(Request $request)
{
$data = $request->post();
$validate = Validate::rule([
'product_id' => 'require|integer',
'quantity' => 'require|integer|between:1,99',
])->check($data);
if (!$validate) {
return json(['code' => 400, 'msg' => '参数错误']);
}
Db::startTrans();
try {
$order = OrderModel::create([
'user_id' => $request->uid,
'product_id' => $data['product_id'],
'quantity' => $data['quantity'],
'amount' => $this->calcAmount($data),
'status' => 'created',
]);
// 扣库存
ProductModel::where('id', $data['product_id'])
->dec('stock', $data['quantity'])
->update();
Db::commit();
} catch (Throwable $e) {
Db::rollback();
throw $e;
}
// 下面这些都是副作用
try {
SmsService::send($order->user_id, '下单成功:' . $order->order_no);
} catch (Throwable $e) {
trace('短信发送失败:' . $e->getMessage(), 'error');
}
try {
PointService::grant($order->user_id, intval($order->amount * 10));
} catch (Throwable $e) {
trace('积分发放失败:' . $e->getMessage(), 'error');
}
try {
WarehouseService::notify($order->toArray());
} catch (Throwable $e) {
trace('仓库通知失败:' . $e->getMessage(), 'error');
}
try {
StatModel::accumulate($order->created_at, $order->amount);
} catch (Throwable $e) {
trace('统计累计失败:' . $e->getMessage(), 'error');
}
try {
RiskService::snapshot($order->user_id, $order->toArray());
} catch (Throwable $e) {
trace('风控快照失败:' . $e->getMessage(), 'error');
}
return json(['code' => 0, 'data' => ['order_no' => $order->order_no]]);
}
}
读一遍就能看出几个问题。每一个副作用都包着 try-catch,把异常吞掉只记日志。这是不得已——因为如果任何一个副作用抛异常,用户就会看到”下单失败”,但实际上订单已经写进去了。用户重试就会重复下单。
而吞掉异常也带来新问题:短信没发出去,用户没收到,客服也不知道。日志里躺着一条 error,一周之后翻出来才发现。
另外,五个 try-catch 加起来占了八十行,真正和”下单”这件事相关的业务逻辑被淹没了。以后每次要加新副作用,就往这里塞一段;要改某个副作用的逻辑,得先找到它在第几段里。
二、第一次重构:把副作用写成监听器
ThinkPHP 8 的事件系统用起来很轻。先定义一个事件类:
namespace appevent;
class OrderPlaced
{
public function __construct(
public readonly int $orderId,
public readonly string $orderNo,
public readonly int $userId,
public readonly string $amount,
public readonly string $placedAt,
) {}
}
事件类只承载数据,逻辑放在监听器里。这是 ThinkPHP 事件系统的一个约定——事件类不做事,它是”发生了什么”的快照。
然后在配置文件里注册监听关系:
// app/event.php
return [
'bind' => [
'OrderPlaced' => appeventOrderPlaced::class,
],
'listen' => [
'OrderPlaced' => [
applistenerorderSendSms::class,
applistenerorderGrantPoints::class,
applistenerorderNotifyWarehouse::class,
applistenerorderAccumulateStat::class,
applistenerorderRiskSnapshot::class,
],
],
];
每个监听器都很短:
namespace applistenerorder;
use appeventOrderPlaced;
use appserviceSmsService;
class SendSms
{
public function handle(OrderPlaced $event)
{
SmsService::send(
$event->userId,
'下单成功:' . $event->orderNo
);
}
}
控制器里下单逻辑保持原样,只在事务提交之后加一行:
Db::commit();
event('OrderPlaced', new appeventOrderPlaced(
orderId: $order->id,
orderNo: $order->order_no,
userId: $order->user_id,
amount: $order->amount,
placedAt: $order->created_at,
));
return json(['code' => 0, 'data' => ['order_no' => $order->order_no]]);
到这里,控制器瘦回来了。但有个新问题:这五个监听器里任何一个抛异常,都会中断后面监听器的执行。因为默认情况下 ThinkPHP 的事件派发是同步的、顺序执行的。
换句话说,从”五个 try-catch 各自吞异常”变成了”一次异常全部扑街”,稳定性的方向走反了。
三、第二次重构:用订阅者分组,并且隔离异常
先解决分组问题。event.php 里把五个监听器全列出来,随着项目变大,这个文件会变成一个大杂烩。ThinkPHP 提供了订阅者机制来自我注册:
namespace appsubscribe;
use thinkEvent;
class OrderSubscribe
{
public function subscribe(Event $event): void
{
$event->listen('OrderPlaced', [
applistenerorderSendSms::class,
applistenerorderGrantPoints::class,
applistenerorderNotifyWarehouse::class,
applistenerorderAccumulateStat::class,
applistenerorderRiskSnapshot::class,
]);
$event->listen('OrderPaid', [
applistenerorderSendPaidNotice::class,
]);
$event->listen('OrderRefunded', [
applistenerorderRollbackStock::class,
applistenerorderDeductPoints::class,
]);
}
}
然后在 app/event.php 里注册订阅者:
return [
'subscribe' => [
appsubscribeOrderSubscribe::class,
],
];
这样订单相关的所有事件都集中在一个订阅者文件里,找起来方便,加新事件也不需要来回跳。
但分组没有解决异常隔离。要处理这个,最简单的方式是在每个监听器里自己包 try-catch。不过写多了会累,可以在监听器基类里做一层包装:
namespace applistener;
abstract class BaseListener
{
final public function handle($event)
{
try {
$this->process($event);
} catch (Throwable $e) {
// 记住是哪一类事件、哪个监听器失败
trace(sprintf(
'[listener] %s 处理 %s 失败:%s',
static::class,
get_class($event),
$e->getMessage()
), 'error');
}
}
abstract protected function process($event): void;
}
所有监听器都 extends BaseListener 并实现 process,这样异常隔离就是默认行为,不会因为某个监听器忘了写 try-catch 而破坏整个派发流程。
当然这也带来一个副作用:异常被吞了。如果某些监听器的失败必须显式通知某个地方,就得重写 handle 方法,或者在 process 里自己处理。这是必要的灵活性。
四、第三次重构:把能异步的都放到队列里
监听器分离之后,下单接口的执行时间反而长了。因为现在同步执行的监听器多了五个,每个哪怕只花 20 毫秒,总共也拖了 100 毫秒。用户点击下单之后要多等这么久才看到结果。
但仔细想想,这五个副作用里,只有库存扣减是用户在意的。短信可以晚一点发,积分晚一秒到也行,仓库通知和统计更是没人着急。真正需要”立刻完成”的只有订单写入和库存扣减。
所以下一步是把大部分监听器改成异步入队。
装 think-queue
composer require topthink/think-queue
然后配置驱动。生产环境用 redis,本地开发可以用 database 甚至 sync:
// config/queue.php
return [
'default' => env('QUEUE_DRIVER', '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' => 'database',
'table' => 'failed_jobs',
],
];
本地用 sync 是很舒服的一件事:队列任务直接同步执行,不需要额外启 worker,debug 起来和普通函数没区别。上线切到 redis 就行,业务代码完全不用改。
把监听器改成入队
namespace applistenerorder;
use appeventOrderPlaced;
use thinkfacadeQueue;
class SendSms extends applistenerBaseListener
{
protected function process($event): void
{
Queue::push(
appjobSendOrderSms::class,
[
'userId' => $event->userId,
'orderNo' => $event->orderNo,
],
'notify' // 队列名
);
}
}
监听器本身还是同步执行的,但它做的事情从”发短信”变成了”入队”。入队操作很快,几毫秒就够了,对下单接口的响应时间影响可以忽略。
Job 类里才是真正的业务:
namespace appjob;
use appserviceSmsService;
use thinkqueueJob;
class SendOrderSms
{
public function fire(Job $job, array $data): void
{
try {
SmsService::send(
$data['userId'],
'下单成功:' . $data['orderNo']
);
$job->delete();
} catch (Throwable $e) {
if ($job->attempts() >= 3) {
// 超过重试次数,放弃并落库,人工介入
$job->delete();
thinkfacadeLog::error('短信任务最终失败', [
'data' => $data,
'error' => $e->getMessage(),
]);
return;
}
// 延迟 30 秒重试
$job->release(30);
}
}
}
这段代码大概解释了 think-queue 的三个核心动作:delete() 表示任务完成,从队列移除;release($delay) 表示任务失败要重试,把消息放回队列并延迟执行;attempts() 返回当前是第几次尝试。
启动 worker:
php think queue:work --queue notify --daemon
或者用 supervisor 守护进程。生产环境基本不会用 --daemon,因为 worker 内存泄漏问题会在长时间运行后显现,交给 supervisor 定时重启更稳。
五、幂等:队列任务绕不开的一道坎
异步之后有个新问题:任务可能被消费多次。
Redis 队列在网络抖动、worker 崩溃、release 超时等情况下,同一任务被投递两次以上是常见现象。所以每个 Job 都要考虑”如果这个任务已经被执行过了,怎么办”。
对于”发短信”这个任务,用户收到两条相同内容可能还能忍。但对于”发积分”、”扣余额”这种,重复执行就是事故。
通用做法是用一个唯一键加锁:
namespace appjob;
use thinkfacadeCache;
class GrantPoints
{
public function fire(Job $job, array $data): void
{
$lockKey = 'job:grant_points:order:' . $data['orderNo'];
// Redis 的 SETNX + EX,5 分钟过期
$locked = Cache::store('redis')->handler()->set(
$lockKey,
1,
['nx', 'ex' => 300]
);
if (!$locked) {
// 已经有别的 worker 在跑,或已经跑过,直接丢弃
$job->delete();
return;
}
try {
appservicePointService::grant(
$data['userId'],
(int) $data['points']
);
$job->delete();
} catch (Throwable $e) {
// 任务失败,解锁以便下次重试
Cache::store('redis')->delete($lockKey);
$job->attempts() >= 3 ? $job->delete() : $job->release(30);
}
}
}
这里锁的语义跟”任务执行中”和”任务已完成”共用了一个键,有优点也有缺点。优点是逻辑简单,重试时锁已经释放(因为失败路径主动删了锁)。缺点是任务正常完成后锁要等到 5 分钟过期,这段时间内如果有重复投递会被识别为已执行——正好是我们要的效果。
但如果业务上要求”允许用户 5 分钟之后再对同一个订单重新发起某项操作”,就得更精细地设计 key,比如加上”任务类型 + 订单号 + 状态标志”。这属于具体情况具体分析,不是队列本身的事。
六、六个实际踩过的坑
1. 事务未提交就派发事件
这是最容易翻车的一个。
Db::startTrans();
try {
$order = OrderModel::create([...]);
// 错误:此时事务还没提交
event('OrderPlaced', new OrderPlaced(...));
Db::commit();
} catch (Throwable $e) {
Db::rollback();
throw $e;
}
如果监听器里查这张表、或者查询关联表,在 MySQL 默认隔离级别下可能看不到刚插入的数据(虽然同一个连接里通常可以,但其他连接就看不到)。而如果监听器里做了写操作,事务回滚时这些操作不会跟着回滚。更严重的场景是监听器里直接把订单数据推给下游——下游拿到的是一个可能根本不存在的订单。
正确写法:先提交,再派发。
Db::startTrans();
try {
$order = OrderModel::create([...]);
Db::commit();
} catch (Throwable $e) {
Db::rollback();
throw $e;
}
// 事务已经提交
event('OrderPlaced', new OrderPlaced(
orderId: $order->id,
// ...
));
这一条看似很简单,但几乎每个用 ThinkPHP 事件的项目都会犯一次。写之前先想一想”我这个事件要在事务内还是事务外派发”。
2. Job 类名变更会破坏老数据
队列里的消息是 JSON 序列化之后的 [类名, 参数]。如果你改了 Job 类的命名空间,比如从 appjobSendSms 挪到了 appjoborderSendSms,那些还在队列里的老消息在反序列化时会失败。
think-queue 在找不到类的时候不会报错退出,而是静默丢弃。所以线上会出现”改完之后某段时间有些任务莫名消失”的现象,排查起来非常费劲。
几种处理方式:
最好的是保留一个兼容类,转发到新类:
namespace appjob;
// 已废弃:只为了兼容老队列,新代码不要用
class SendSms
{
public function fire($job, $data)
{
(new appjoborderSendSms())->fire($job, $data);
}
}
或者部署前先停 worker,等队列清空再上。这个更激进,但更干净。
3. attempts 判断写错导致无限重试
if ($job->attempts() >= 3) {
$job->delete();
return;
}
$job->release(30);
看着没问题。但如果写成这样:
if ($job->attempts() > 3) {
// ...
}
那就允许了 4 次尝试,而不是 3 次。attempts() 从 1 开始计数,任务第一次进入 fire 方法时它就返回 1。习惯从 0 开始计数的开发者容易写错。
更麻烦的一种写错是:只在 catch 里判断 attempts,而 release 又放在 catch 外面。这样第一次失败的 release 被调用了,第二次失败的 release 又调用了,attempts 一直没超过 3 是因为任务本身可能有变化。这种情况很罕见,但如果业务代码里有幂等或者在 catch 里 return 了没走 release,就容易写乱。
简单原则:把 attempts 判断和 release 放进同一个分支里,两者不要分开。
4. Job 参数里塞了无法序列化的对象
Queue::push(SendSms::class, [
'user' => $userModel, // Model 对象,序列化后是一堆属性
'order' => $order, // 同上
'cfg' => config('sms'), // 数组,还好
'fn' => fn($x) => $x * 2, // 闭包,直接报错
'fp' => fopen('/tmp/x', 'r'), // resource,报错
]);
Job 参数走的是 JSON 序列化。Model 对象序列化之后会变成属性数组,读回来的时候是一个 stdClass 或者 Model(取决于配置),行为可能和原对象不同。闭包和资源直接会抛异常。
最简单也最稳的做法是:只传标量和数组。
Queue::push(SendSms::class, [
'userId' => $user->id,
'orderId' => $order->id,
'templateCode' => 'ORDER_PLACED',
]);
Job 内部再按 id 去查需要的信息。虽然多一次查询,但消息体积小、兼容性好、也不会因为 Model 结构变化而让老消息解析失败。
5. queue:work 内存持续增长
worker 是常驻进程。处理几万个任务之后,内存会明显上涨。这不是 think-queue 的问题,是 PHP 本身在某些场景(大数组、循环引用)下内存回收不及时。
think-queue 提供了内存上限参数:
php think queue:work --queue notify --daemon --memory 256
当 worker 使用的内存超过 256MB 时,它会处理完当前任务之后退出。supervisor 会把它重新拉起来。
容器环境下,如果不想跑 supervisor,也可以用 restart 命令定期重启所有 worker:
php think queue:restart
这个命令会在 Redis 里写一个时间戳,所有 worker 发现这个时间戳变了就会在下一次循环时退出。放到 crontab 里每小时跑一次就行。
6. 改了 Job 代码但 worker 不生效
这个坑第一次遇到的时候很迷惑。明明改了 Job 类的代码,测试环境一跑就有效,生产环境怎么都不生效。
原因:生产环境的 worker 已经启动了很久,PHP 已经把 Job 类加载到内存里了。你部署新代码之后,worker 进程还在跑老的代码,不会自动重新加载。
解决办法就是前面那个 queue:restart。养成习惯:每次部署之后都执行一次。而且这个命令执行得很快,不需要等到 worker 处理完当前任务——它只是让 worker 在下次拿任务前主动退出。
更规范的做法是把 queue:restart 写进部署脚本,最后一步自动触发。
七、事件系统适用与不适用的边界
重构完之后,下单接口从两百多行变成了四十多行,五个副作用的实现分散到了独立的监听器和 Job 里。但事件系统不是万能药,用之前要判断清楚。
适合的场景:一件事发生之后触发多个松耦合的后续动作;”主流程”和”附加流程”分界明显,附加流程失败不应该影响主流程;跨模块通信,模块之间不应该直接引用对方的 Service 类。
不适合的场景:只有一两个副作用,抽成事件反而多了一层间接;副作用的结果影响主流程的返回(比如下单成功之后要返回给用户一个积分 ID)——这种情况仍然应该同步获取结果,事件不适合;需要在事务内保证一致性的操作——事件天生不适合做事务保证,需要另想办法(比如事务消息或 outbox 模式)。
还有一条经验:事件名要反映”已经发生的事实”,不要用”命令”的语气。用 OrderPlaced 而不是 CreateOrder,用 UserRegistered 而不是 RegisterUser。这不是命名洁癖,而是因为事件是过去时,它承载的是”数据已经被写入”这个事实;如果写成命令语气,容易让调用方误以为可以在事件里做”写数据”的动作,把事件当成了另一种形式的服务调用。
八、写在最后
这套重构从提出到上线大概花了两周,中间讨论最多的不是技术实现,而是”哪几个副作用可以异步”这个业务判断。
比如积分发放,一开始定的是异步。上线之后发现用户下单成功之后立刻去积分页面看不到刚到账的积分,客服被打了好几个电话。后来改成同步调 Service 层,但把点数计算和写库分离,实际还是异步写。技术方案的正确性是一回事,用户体验上的”可控延迟”是另一回事。
还有一件事想强调一下:事件系统的核心价值不在”少写代码”,而在于把业务模块之间的关系显式化。以前订单模块直接 import 了短信模块的类,这种跨模块依赖被隐藏在 import 语句里,看代码的时候很难评估订单模块到底依赖了多少东西。事件系统把这种依赖变成了一张清单,就在 OrderSubscribe 里,一眼看得清清楚楚。
对一个小项目,这可能不算什么;但对一个功能越加越多的中大型项目,这张清单本身就是很好的一份架构文档。

