Laravel队列实战:从配置到幂等设计,避免重复消费
2026/9/9 17:56:37 网站建设 项目流程

做了三年多的 Laravel 项目,我几乎每个系统里都会用到消息队列,但真正让我把这一块彻底吃透的,是去年接手的一个订单通知加数据同步的项目。那段时间每天都要面对队列堆积、任务超时、重复消费这类问题,光是排查 worker 挂掉的原因就折腾了好几个通宵。这篇文章把我这段时间积累的 Laravel 队列实操经验完整梳理一遍,从配置、Job 编写到幂等设计、结果存储,全部是我实际跑过生产环境后验证过的方案,适合正在用 Laravel 做中后台系统、希望把队列用得更稳的开发者。


1. 内容整体设计与思路拆解

1.1 为什么业务一复杂就必须上消息队列

先说一个很典型的场景。用户注册成功后,系统要发欢迎邮件、发短信验证码、给推荐人加积分、往数据分析平台推送用户事件。如果你把这些逻辑全部写在注册控制器里同步执行,一次注册请求的耗时可能从 50 毫秒直接飙到 2 秒以上,而且任何一个第三方接口超时,用户就得一直转圈等待。更麻烦的是,邮件服务商宕机、短信通道限流的时候,整个注册接口都会跟着报错。

消息队列解决的就是这个耦合问题。它的核心思路很简单:把耗时操作、非核心操作从请求链路中剥离出来,先返回“已受理”,再通过后台 worker 进程逐步处理。这样用户感受到的响应速度变快了,核心业务也不容易被外部依赖拖垮。我经常用一个类比:同步请求就像去餐厅点餐后站在厨师旁边等菜,消息队列则是拿了号牌回座位等,菜做好了服务员自然会端上来。

Laravel 的队列系统把这一整套机制做成了标准组件,你不需要自己去实现 Redis 的 BRPOPLPUSH、不需要手动管理 worker 进程,只需要定义好 Job、调用 dispatch、跑一个 queue:work 命令,剩下的事框架都帮你兜住了。但正因为它“太方便”,很多人反而忽略了背后的一些设计细节,导致业务量一上来就出问题。

1.2 Laravel 队列的架构选型:驱动、进程模型与 Job 生命周期

Laravel 队列从架构上可以拆成四层:驱动层、队列层、Job 层、Worker 层。

驱动层支持 database、redis、sqs、sync 等多种方案。我个人最推荐 Redis,原因很简单:Redis 的 BRPOPLPUSH 和 Stream 数据结构原生态支持阻塞读取和消息确认,性能远超数据库轮询,可靠性又比纯内存队列好得多。数据库驱动更适合学习和小流量项目,但生产环境一旦并发上来,每秒几千次查询去轮询 jobs 表,数据库很快就扛不住了。

进程模型大概是这样的:queue:work启动一个常驻 PHP 进程,这个进程在 while 循环里不断从队列中 pop 任务,然后执行。PHP 的常驻进程和传统 PHP-FPM 不同,同一进程会连续处理多个任务,所以内存泄漏、全局状态污染这些问题在长时间运行后可能会浮现出来,这也是为什么 Laravel 官方推荐用 supervisor 监控 worker 进程、在达到内存上限时自动重启。

一个 Job 的完整生命周期是:客户端调用 dispatch 将任务推送到队列 -> Redis 中存储序列化后的任务数据 -> worker 从队列中取出任务 -> 框架调用 Job 的 handle 方法 -> 执行成功或抛异常 -> 若抛异常则根据重试策略决定马上重试、延迟重试还是进失败表。理解这个生命周期非常重要,因为后面很多疑难杂症,比如任务重复执行、任务丢失、超时被 kill,全都出在这些环节的衔接处。


2. 环境准备与核心配置实操

2.1 队列配置逐项拆解:从 queue.php 到 .env

把队列用起来的第一步不是写 Job,而是把配置搞清楚。我见过太多的团队在.env里写了QUEUE_CONNECTION=redis之后就不管了,结果任务死活不执行,最后发现是 Redis 连接配置不对或者队列名没对应上。

Laravel 的队列配置在config/queue.php,其中default键对应.env中的QUEUE_CONNECTION。如果你用 Redis 驱动,需要关注这几个配置项:

'redis' => [ 'driver' => 'redis', 'connection' => 'default', 'queue' => env('REDIS_QUEUE', 'default'), 'retry_after' => 90, 'block_for' => 5, 'after_commit' => false, ],

retry_after这个参数很多人会忽略,但它极其重要。它的意思是:如果 worker 处理任务超过了这个秒数,框架就认为任务处理失败了,会让任务重新回到队列中可被其他 worker 领取。注意这里有个坑:retry_after必须大于你单个任务处理的最长耗时,否则一个正常运行但耗时较久的任务就会被重复领取执行,引发重复消费问题。

block_for是 worker 阻塞等待任务的时间,默认 5 秒。这个参数影响的是 Redis 驱动的 BRPOP 行为,设置为 5 表示没有任务时 worker 会阻塞 5 秒再去轮询,可以有效降低 Redis 的压力。

.env里常用配置示例:

QUEUE_CONNECTION=redis REDIS_HOST=127.0.0.1 REDIS_PORT=6379 REDIS_PASSWORD=null REDIS_QUEUE=default

2.2 从 sync 切换到 Redis:别在生产环境踩这个坑

很多新项目开发时习惯用默认的QUEUE_CONNECTION=sync,这意味着任务会被立即同步执行,方便调试。但如果你把sync模式的项目直接部署到线上,问题就来了:队列任务不仅没有异步化,而且所有耗时操作都堆在请求进程里执行,性能和并发能力都会出问题。

我踩过一次很深刻的坑:有个项目在本地用sync模式调试发送邮件,一切正常,部署到测试环境也改了QUEUE_CONNECTION=redis,但因为没有重启queue:work进程,旧的 worker 还保持着旧的配置,任务一直在老进程里跑。最终排查的时候发现,队列任务根本没走 Redis,全是旧 worker 进程用数据库连接处理的。所以每次修改队列配置,一定要记得重启 worker:

php artisan queue:restart

这条命令会给所有 worker 发信号,让它们在处理完当前任务后优雅退出,然后由 supervisor 或你手动重启。

另外,after_commit这个配置项也很值得关注。它控制的是:如果 Job 在数据库事务中分发,是否等事务提交后再真正把任务推入队列。默认是 false,意味着事务还没提交,任务就已经进了队列,而 worker 如果立刻去读数据库,可能读到的是旧数据。这在数据一致性要求高的场景下是个隐患。我的建议是,涉及事务操作的场景,把after_commit设为 true,配合dispatch函数使用:

DB::transaction(function () { $order = Order::create([...]); DispatchOrderNotification::dispatch($order); // 事务提交后才真正入队 });

3. Job 任务的创建、分发与执行链路

3.1 创建 Job 与分发方式:dispatch、队列选择与延迟执行

生成一个 Job 很简单:

php artisan make:job ProcessVideo

生成的 Job 类位于app/Jobs目录,默认包含一个handle方法。最重要的区分是:构造函数中的参数会被序列化存入队列,handle方法中的参数则是从容器中解析的依赖。

举个例子,你要处理一个视频转码任务,需要传入视频 ID,还要在触发时注入视频处理服务:

class ProcessVideo implements ShouldQueue { public $videoId; public $timeout = 300; public function __construct($videoId) { $this->videoId = $videoId; } public function handle(VideoProcessor $processor) { $processor->process($this->videoId); } }

分发方式有几种,实际使用中我比较习惯用 dispatch 辅助函数和 Job 的静态方法:

ProcessVideo::dispatch($videoId); ProcessVideo::dispatch($videoId)->onQueue('high'); // 指定队列 ProcessVideo::dispatch($videoId)->onConnection('redis'); // 指定连接 ProcessVideo::dispatch($videoId)->delay(now()->addMinutes(10)); // 延迟执行

onQueue是非常实用的功能。我通常把队列按优先级拆成多个:high队列处理注册邮件、支付回调这类时效性强的任务;default队列处理普通的通知;low队列处理数据统计、报表导出这类不着急的任务。启动 worker 的时候用逗号分隔:

php artisan queue:work redis --queue=high,default,low --tries=3 --timeout=120

这样 worker 会优先消费high队列的所有任务,再按顺序处理后面的队列。

3.2 消费端的关键参数:tries、backoff、timeout、retryUntil 怎么配

很多人的队列任务出问题,就是因为这几个参数没搞懂,或者说没搞懂之间的关系。我一个个说:

tries表示任务最多执行几次。如果超过这个次数仍然抛异常,任务会进入failed_jobs表。在 Job 类中可以直接定义:

public $tries = 3;

backoff表示每次重试之间的等待时间。可以是整数秒,也可以是数组,表示不同次数的不同等待时间:

public $backoff = [10, 60, 300];

上面这个数组表示:第一次失败后等 10 秒重试,第二次失败后等 60 秒,第三次失败后等 300 秒(如果还是有 try 的话)。

timeout表示单个任务允许执行的最大秒数。超过这个时间,worker 会抛出一个MaxAttemptsExceededException,任务被终止。必须注意,timeout要小于retry_after,否则 worker 进程被 kill 掉之后,任务在retry_after到期前不会被重新领取,会卡在“处理中”状态很长时间。

retryUntiltries是两种不同的重试终止策略。retryUntil指定一个过期时刻,只要没到这个时间点,任务就会一直重试,适合那种“今天之内必须确保发出去”的短信、邮件场景:

public function retryUntil() { return now()->addHours(12); }

我个人的配置经验是:外部 API 调用类任务,用tries=3+backoff=[5,30,120]+timeout=30;本地数据处理类任务,用tries=5+timeout=300;发送通知类任务,用retryUntil保证时效性。

3.3 实战:用 orderBy 和 groupBy 取最新一条且去重的数据

在队列任务里,经常要从一批数据中取每个分类或每个用户的“最新一条”。比如统计用户最新一次登录时间、给每个用户生成最新一条动态的通知等。很多同事写 SQL 的时候习惯直接orderBygroupBy一起用,结果取出来的数据根本不是想要的。

直接写成这样是错的:

// 错误写法:分组后取到的不是每个分组内最新的那一条 $latestMessages = DB::table('messages') ->select('user_id', 'content', 'created_at') ->groupBy('user_id') ->orderBy('created_at', 'desc') ->get();

这个 SQL 在 MySQL 中虽然能执行(只开了 ONLY_FULL_GROUP_BY 会直接报错),但contentcreated_at并不是每个用户最新一条记录的内容,而是分组后任意选中的一条,只碰巧看起来像。原因在于:groupBy先按user_id分组,然后orderBy是在分组结果上排序,分组内的其他列怎么选取,MySQL 不保证。

正确做法之一:先用子查询配合groupBy拿到每个user_id的最大id,再用id去关联原表取完整数据:

$subQuery = DB::table('messages') ->select('user_id', DB::raw('MAX(id) as max_id')) ->groupBy('user_id'); $latestMessages = DB::table('messages as m') ->joinSub($subQuery, 'latest', function ($join) { $join->on('m.id', '=', 'latest.max_id'); }) ->select('m.user_id', 'm.content', 'm.created_at') ->get();

还有一种更通用的写法是用窗口函数,MySQL 8.0+ 支持ROW_NUMBER()

$latestMessages = DB::select(" SELECT user_id, content, created_at FROM ( SELECT user_id, content, created_at, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY created_at DESC) AS rn FROM messages ) t WHERE t.rn = 1 ");

在队列任务的实际场景里,这样取到最新数据之后,再进行发送通知、聚合统计等后续处理,就能避免任务重复消费后把旧数据也带出来。当你用tries=3去重试一个任务时,如果不做数据过滤,很可能第二次执行时会把同一批数据全量处理一遍,造成重复通知。这里的去重和取最新数据,实际上也是在给任务消费的数据做了一层“幂等控制”。


4. 消息队列重复消费问题:从根源到幂等方案

4.1 重复消费为什么会发生:worker 崩溃、超时与网络波动

消息队列领域有句著名的话:队列只保证 at-least-once,不保证 exactly-once。也就是说,同一个任务被重复执行是常态,不是意外。Laravel 队列也不例外。

重复消费的根源主要有三个。第一,worker 在处理任务过程中崩溃(比如被 OOM kill、部署时被杀掉、PHP 致命错误),任务已经执行了一部分但没有返回成功,Redis 中这个任务还在 pending 状态,retry_after到期后会被其他 worker 重新领取。第二,timeout设置得太小,任务还在慢速处理中就被 worker 判定超时,然后重新入队。第三,使用 Redis Stream 作为驱动时,消费者组 ACK 失败、消费进程在读取后未确认就宕机,消息会在 PEL(Pending Entries List)中留存,重新消费时再次投递。

我举一个实际业务案例。有一次我们在队列任务里给用户发优惠券,任务内容是“读取用户记录 -> 生成券码 -> 发到用户账户”。某次 Redis 高峰期,任务执行到“生成券码”之后、还没走到“发到用户账户”,worker 被 supervisord 因内存超限重启了。retry_after90 秒到期后,同一个任务被另一个 worker 领取,再次从“读取用户记录”开始执行,于是这个用户收到了两张相同的优惠券。

事后我的总结是:与其想办法让队列系统不重复投递,不如接受这个现实,在业务层做幂等处理。

4.2 幂等设计的三种落地模式

幂等设计的核心只有一个:让同一个任务无论执行多少次,最终产生的业务结果都一样。结合 Laravel 的实践,我常用三种模式。

第一种是唯一约束模式。适用于“同一个业务只能有一条记录”的场景。比如上面的发券案例,在user_coupons表加一个(user_id, coupon_template_id)唯一索引,任务执行时用firstOrCreate或者 insertOrIgnore:

UserCoupon::firstOrCreate([ 'user_id' => $this->userId, 'coupon_template_id' => $this->templateId, ], [ 'coupon_code' => $couponCode, 'status' => 'active', ]);

这样就算任务被重复执行,第二次也会因为唯一索引冲突而直接跳过。

第二种是 Redis 标记模式。适用于“处理动作不可用唯一约束表达”的场景。比如任务要调用第三方短信 API 发送一条验证码,无法通过在本地表里加索引实现幂等,那就在执行前用一个业务标识(如订单号、用户 ID 加任务类型)写入 Redis,设置一个合理的过期时间:

$lockKey = 'sms:send:' . $this->userId . ':' . $this->scene; $locked = Redis::set($lockKey, '1', 'EX', 300, 'NX'); if (!$locked) { return; // 说明已经处理过了 } try { $smsService->send($this->userId, $this->content); } catch (\Throwable $e) { Redis::del($lockKey); // 失败时释放锁,允许重试 throw $e; }

这里用NX参数保证原子性,避免并发时两个进程同时拿到锁。

第三种是状态机模式。适用于“业务有明确状态流转”的场景。比如订单状态从pendingpaid,只能转换一次。任务执行前先检查当前状态是否已经是目标状态,是则直接返回:

$order = Order::find($this->orderId); if ($order->status === 'paid') { return; // 已处理,直接跳过 }

这三种模式可以组合使用。我在实际项目中通常是“唯一索引兜底 + Redis 标记防止重复调用外部接口 + 状态机保护核心业务流转”,三层一起上,才能保证队列任务在生产环境长时间运行不出问题。

4.3 参考 broker + backend 双存储设计:Redis 同时做任务队列和结果存储

很多人在用 Laravel 队列时只关心任务能不能被消费,却不关心任务执行的结果。实际上,在微服务和异步任务系统中,任务结果的可追溯性非常重要。Celery 的架构里有一个“broker + backend”双存储模式:broker 存储任务消息本身,backend 存储任务执行结果。Laravel 虽然没有内置完全等价的机制,但用 Redis 完全可以实现一套类似的设计。

具体做法是:让 Laravel 队列使用 Redis 作为 broker(这是框架原生支持的),然后在 Job 内部把执行结果写入另一组独立的 Redis key,作为 backend。这样任务执行完,你不需要翻日志就能查到每个任务跑完后的状态、返回值和耗时。

我封装过一个简单的 Job 基类:

abstract class TrackableJob implements ShouldQueue { public $timeout = 60; abstract public function execute(); public function handle() { $jobId = $this->job->getJobId(); $resultKey = 'job:result:' . $jobId; Redis::hset($resultKey, 'status', 'processing'); Redis::hset($resultKey, 'started_at', now()->toDateTimeString()); Redis::expire($resultKey, 86400); // 保留24小时 try { $result = $this->execute(); Redis::hset($resultKey, [ 'status' => 'success', 'result' => json_encode($result), 'finished_at' => now()->toDateTimeString(), ]); } catch (\Throwable $e) { Redis::hset($resultKey, [ 'status' => 'failed', 'error' => $e->getMessage(), 'finished_at' => now()->toDateTimeString(), ]); throw $e; // 继续走 Laravel 的重试/失败流程 } } }

然后通过一个全局唯一的业务 ID 把任务 ID 和业务关联起来。比如生成报表的任务:

class GenerateReport extends TrackableJob { public $reportId; public function __construct($reportId) { $this->reportId = $reportId; } public function execute() { $data = ReportService::generate($this->reportId); Redis::set('report:result:' . $this->reportId, json_encode($data), 'EX', 86400); return $data; } }

这样你可以在后端提供一个查询接口,前端轮询任务状态,拿到结果后展示,不需要再用一个单独的 timer 去数据库里反复查状态。这种设计在处理耗时较长的异步任务(如导出大批量数据、批量人脸比对等)时特别好用。


5. 常见问题与排查技巧实录

5.1 任务一直 pending 不执行:先查这几个地方

队列任务最常见的现象是:dispatch 之后数据进 Redis 了,但任务就是不跑。我排查这类问题的顺序是固定的。

先确认 worker 是否在运行:

ps aux | grep "queue:work"

如果 worker 没跑,检查 supervisor 配置是否正确。如果 worker 在跑,接着看它监听的是哪个队列。很多时候你 dispatch 任务时用了->onQueue('high'),但启动 worker 时只写了--queue=default,任务就一直留在 Redis 的queues:highkey 里没人消费。用 Redis 客户端查看:

redis-cli llen queues:high

如果长度一直在增长,而queues:default是空的,那基本就是队列名不匹配。还要确认.env里的QUEUE_CONNECTION=redis是否真的生效了。有一个小技巧,在config/queue.php中临时dd(config('queue.default')),看输出值就知道实际走的是哪个驱动。

另外,queue:workqueue:listen的行为有差异,queue:listen每次都会启动一个新的框架实例来处理任务,比较慢;queue:work是常驻进程,性能好很多,生产环境推荐用queue:work

5.2 失败任务与重试机制管理:failed_jobs 表与手动重跑

任务多次执行仍然失败后,会进入failed_jobs表。前提是你已经创建了这张表:

php artisan queue:failed-table php artisan migrate

然后在.env中确认QUEUE_FAILED_DRIVER=database(Laravel 11 及以上是FAILED_JOB_DRIVER=database)。之后查看失败任务:

php artisan queue:failed

你会看到每个失败任务的 ID、Job 类名、失败原因和失败时间。如果确认是外部接口临时故障,可以直接重试:

php artisan queue:retry 5 # 按 ID 重试 php artisan queue:retry all # 重试全部

不需要的任务可以删除:

php artisan queue:forget 5

清空整个失败表:

php artisan queue:flush

这里分享一个经验:不要在failed_jobs表里躺了成千上万条数据之后才想起处理。我通常在 Job 里加一个失败通知:

public function failed(\Throwable $e) { Log::error('队列任务执行失败', [ 'job' => static::class, 'error' => $e->getMessage(), 'data' => $this->payload ?? null, ]); // 发送告警到钉钉/飞书/企业微信 Alert::send('队列任务失败:' . static::class . ',原因:' . $e->getMessage()); }

这样任务一失败就能第一时间知道,而不是等到业务方反馈才去翻日志。

5.3 队列性能优化:并发、批量处理与资源控制

队列性能优化有几个方向。

第一个是并发控制。单台机器的 worker 并不是越多越好,因为每个 worker 都是一个常驻 PHP 进程,占用的内存和 CPU 都不低。如果单个 worker 处理的任务里有大量网络 IO(比如调用外部 API),可以提高 worker 数量,因为进程在等待 IO 时 CPU 是空闲的;如果任务是 CPU 密集型(比如图片处理、数据计算),worker 数量建议接近 CPU 核心数,避免过多的上下文切换。

我常驻服务器的典型配置是 8 核 16G,Redis 和 web 服务都在同一台机器上,队列 worker 我就开 4 个:

[program:laravel-worker] process_name=%(program_name)s_%(process_num)02d command=php /var/www/html/artisan queue:work redis --queue=high,default --tries=3 --timeout=60 numprocs=4 autostart=true autorestart=true stopwaitsecs=3600

stopwaitsecs很关键。Supervisor 在重启 worker 时,会等待 worker 优雅退出,但这个时间不能短于任务最长执行时间,否则 worker 在收到停止信号后还在处理任务,直接被 kill 掉,当前任务就白跑了,还可能引发重复消费。

第二个是批量处理。Laravel 队列默认一条消息一个任务,但对于日志写入、数据统计这类高频低耗任务,可以先把数据攒到内存数组中,满足一定数量后再批量入库。我自己的做法是维护一个BatchLogJob,它接收一个数组,然后一次insert多条日志:

class BatchWriteLog implements ShouldQueue { public $logs; public function __construct(array $logs) { $this->logs = $logs; } public function handle() { LogModel::insert($this->logs); } }

在业务侧,用 Redis 把日志暂时积攒起来,每满 50 条或者每分钟 flush 一次,触发一个 BatchWriteLog:

$key = 'log:buffer'; Redis::rpush($key, json_encode($logData)); $count = Redis::llen($key); if ($count >= 50) { $logs = Redis::lrange($key, 0, -1); Redis::del($key); BatchWriteLog::dispatch(array_map('json_decode', $logs)); }

这样可以把数据库写入次数降到原来的几十分之一。

第三个是失败隔离。大量失败任务如果和正常任务混在同一个队列里,会拖慢正常任务的消费速度。我通常会在failed方法里做一个分类,属于可重试的临时错误就继续抛异常让 Laravel 重试,属于不可恢复的业务错误就直接记录并返回,不再重试:

public function handle() { try { $this->process(); } catch (ExternalServiceException $e) { // 第三方API暂时不可用,抛出异常让框架重试 throw $e; } catch (BusinessLogicException $e) { // 业务上不可能成功(如记录已删除),记录日志后不再重试 Log::warning('任务处理失败,终止重试', ['error' => $e->getMessage()]); } }

5.4 压测时一定要关注的三个指标

队列系统上线前,强烈建议先做一个简单压测,重点关注三个指标:消费速率、堆积积压时间、失败率。

消费速率可以用 Redis 的llen持续观察。如果入队速度大于消费速度,队列长度会持续增长,最终导致任务延迟严重。这时就要考虑增加 worker 数量、升级 Redis 实例、或者优化任务内的业务逻辑(比如把单条插入改成批量插入)。

堆积积压时间指的是从任务入队到被消费的时间差。在handle方法开头加一行日志:

public function handle() { Log::info('任务开始执行', [ 'queue_time' => now()->diffInSeconds($this->createdAt ?? now()), ]); }

当然更优雅的做法是在 Job 构造时记录入队时间,或者利用$this->job->availableAt()来判断。

失败率则是看failed_jobs的增长量。如果某个 Job 的失败率超过 5%,不要简单地增加重试次数,一定要先找出失败的根本原因。我在一个项目中就遇到过一个诡异的现象:某个外部 API 在每天凌晨 2 点准时不可用,重试 3 次还是失败。后来排查发现是对方系统每天凌晨做数据备份,窗口期约 3 分钟。我给对应任务增加了一个retryUntil,让它在这段时间过后继续重试,问题才真正解决。


最后再分享一个我自己的体会:Laravel 队列的入门门槛很低,但用好它需要你对底层机制有足够的敬畏。把retry_aftertimeouttries之间的关系搞清楚,把幂等设计做到位,你的队列系统才能稳定撑住线上流量。如果这篇文章能让你少走几个我走过的弯路,那这篇实操记录就算有价值了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询