ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Laravel队列实战:从配置到幂等设计,避免重复消费

Laravel队列实战:从配置到幂等设计,避免重复消费 做了三年多的 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_CONNECTIONredis之后就不管了结果任务死活不执行最后发现是 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_CONNECTIONredis REDIS_HOST127.0.0.1 REDIS_PORT6379 REDIS_PASSWORDnull REDIS_QUEUEdefault2.2 从 sync 切换到 Redis别在生产环境踩这个坑很多新项目开发时习惯用默认的QUEUE_CONNECTIONsync这意味着任务会被立即同步执行方便调试。但如果你把sync模式的项目直接部署到线上问题就来了队列任务不仅没有异步化而且所有耗时操作都堆在请求进程里执行性能和并发能力都会出问题。我踩过一次很深刻的坑有个项目在本地用sync模式调试发送邮件一切正常部署到测试环境也改了QUEUE_CONNECTIONredis但因为没有重启queue:work进程旧的 worker 还保持着旧的配置任务一直在老进程里跑。最终排查的时候发现队列任务根本没走 Redis全是旧 worker 进程用数据库连接处理的。所以每次修改队列配置一定要记得重启 workerphp 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 --queuehigh,default,low --tries3 --timeout120这样 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到期前不会被重新领取会卡在“处理中”状态很长时间。retryUntil和tries是两种不同的重试终止策略。retryUntil指定一个过期时刻只要没到这个时间点任务就会一直重试适合那种“今天之内必须确保发出去”的短信、邮件场景public function retryUntil() { return now()-addHours(12); }我个人的配置经验是外部 API 调用类任务用tries3backoff[5,30,120]timeout30本地数据处理类任务用tries5timeout300发送通知类任务用retryUntil保证时效性。3.3 实战用 orderBy 和 groupBy 取最新一条且去重的数据在队列任务里经常要从一批数据中取每个分类或每个用户的“最新一条”。比如统计用户最新一次登录时间、给每个用户生成最新一条动态的通知等。很多同事写 SQL 的时候习惯直接orderBy配groupBy一起用结果取出来的数据根本不是想要的。直接写成这样是错的// 错误写法分组后取到的不是每个分组内最新的那一条 $latestMessages DB::table(messages) -select(user_id, content, created_at) -groupBy(user_id) -orderBy(created_at, desc) -get();这个 SQL 在 MySQL 中虽然能执行只开了 ONLY_FULL_GROUP_BY 会直接报错但content和created_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 );在队列任务的实际场景里这样取到最新数据之后再进行发送通知、聚合统计等后续处理就能避免任务重复消费后把旧数据也带出来。当你用tries3去重试一个任务时如果不做数据过滤很可能第二次执行时会把同一批数据全量处理一遍造成重复通知。这里的去重和取最新数据实际上也是在给任务消费的数据做了一层“幂等控制”。4. 消息队列重复消费问题从根源到幂等方案4.1 重复消费为什么会发生worker 崩溃、超时与网络波动消息队列领域有句著名的话队列只保证 at-least-once不保证 exactly-once。也就是说同一个任务被重复执行是常态不是意外。Laravel 队列也不例外。重复消费的根源主要有三个。第一worker 在处理任务过程中崩溃比如被 OOM kill、部署时被杀掉、PHP 致命错误任务已经执行了一部分但没有返回成功Redis 中这个任务还在 pending 状态retry_after到期后会被其他 worker 重新领取。第二timeout设置得太小任务还在慢速处理中就被 worker 判定超时然后重新入队。第三使用 Redis Stream 作为驱动时消费者组 ACK 失败、消费进程在读取后未确认就宕机消息会在 PELPending Entries List中留存重新消费时再次投递。我举一个实际业务案例。有一次我们在队列任务里给用户发优惠券任务内容是“读取用户记录 - 生成券码 - 发到用户账户”。某次 Redis 高峰期任务执行到“生成券码”之后、还没走到“发到用户账户”worker 被 supervisord 因内存超限重启了。retry_after90 秒到期后同一个任务被另一个 worker 领取再次从“读取用户记录”开始执行于是这个用户收到了两张相同的优惠券。事后我的总结是与其想办法让队列系统不重复投递不如接受这个现实在业务层做幂等处理。4.2 幂等设计的三种落地模式幂等设计的核心只有一个让同一个任务无论执行多少次最终产生的业务结果都一样。结合 Laravel 的实践我常用三种模式。第一种是唯一约束模式。适用于“同一个业务只能有一条记录”的场景。比如上面的发券案例在user_coupons表加一个(user_id, coupon_template_id)唯一索引任务执行时用firstOrCreate或者 insertOrIgnoreUserCoupon::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参数保证原子性避免并发时两个进程同时拿到锁。第三种是状态机模式。适用于“业务有明确状态流转”的场景。比如订单状态从pending到paid只能转换一次。任务执行前先检查当前状态是否已经是目标状态是则直接返回$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 时只写了--queuedefault任务就一直留在 Redis 的queues:highkey 里没人消费。用 Redis 客户端查看redis-cli llen queues:high如果长度一直在增长而queues:default是空的那基本就是队列名不匹配。还要确认.env里的QUEUE_CONNECTIONredis是否真的生效了。有一个小技巧在config/queue.php中临时dd(config(queue.default))看输出值就知道实际走的是哪个驱动。另外queue:work和queue:listen的行为有差异queue:listen每次都会启动一个新的框架实例来处理任务比较慢queue:work是常驻进程性能好很多生产环境推荐用queue:work。5.2 失败任务与重试机制管理failed_jobs 表与手动重跑任务多次执行仍然失败后会进入failed_jobs表。前提是你已经创建了这张表php artisan queue:failed-table php artisan migrate然后在.env中确认QUEUE_FAILED_DRIVERdatabaseLaravel 11 及以上是FAILED_JOB_DRIVERdatabase。之后查看失败任务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 核 16GRedis 和 web 服务都在同一台机器上队列 worker 我就开 4 个[program:laravel-worker] process_name%(program_name)s_%(process_num)02d commandphp /var/www/html/artisan queue:work redis --queuehigh,default --tries3 --timeout60 numprocs4 autostarttrue autorestarttrue stopwaitsecs3600stopwaitsecs很关键。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_after、timeout、tries之间的关系搞清楚把幂等设计做到位你的队列系统才能稳定撑住线上流量。如果这篇文章能让你少走几个我走过的弯路那这篇实操记录就算有价值了。
返回列表