队列
简介
在构建 Web 应用程序时,您可能有一些任务(例如解析和存储上传的 CSV 文件)在典型的 Web 请求期间执行耗时过长。幸运的是,Laravel 允许您轻松创建可在后台处理的队列任务。通过将耗时任务转移到队列中,您的应用程序可以以极快的速度响应 Web 请求,并为客户提供更好的用户体验。
Laravel 队列为各种不同的队列后端提供统一的队列 API,例如 Amazon SQS、Redis,甚至是关系型数据库。
Laravel 的队列配置选项存储在应用程序的 config/queue.php 配置文件中。在此文件中,您将找到框架自带的每个队列驱动程序的连接配置,包括 database、Amazon SQS、Redis 和 Beanstalkd 驱动程序,以及一个同步驱动程序(用于开发或测试,可立即执行任务)。框架还包含一个 null 队列驱动程序,用于丢弃已排队的任务。
Laravel Horizon 是一个美观的仪表盘和配置系统,专为 Redis 驱动的队列设计。有关更多信息,请查看完整的 Horizon 文档。
连接与队列
在开始使用 Laravel 队列之前,理解“连接”和“队列”之间的区别非常重要。在 config/queue.php 配置文件中,有一个 connections 配置数组。此选项定义了到 Amazon SQS、Beanstalk 或 Redis 等后端队列服务的连接。然而,任何给定的队列连接可能拥有多个“队列”,这些队列可以被视为队列任务的不同堆栈或堆叠。
请注意,queue 配置文件中的每个连接配置示例都包含一个 queue 属性。这是将任务发送到给定连接时默认分发到的队列。换句话说,如果您在没有明确定义任务应分发到哪个队列的情况下分发任务,该任务将被放入连接配置的 queue 属性中定义的队列。
1use App\Jobs\ProcessPodcast;2 3// This job is sent to the default connection's default queue...4ProcessPodcast::dispatch();5 6// This job is sent to the default connection's "emails" queue...7ProcessPodcast::dispatch()->onQueue('emails');
有些应用程序可能根本不需要将任务推送到多个队列,而是倾向于使用一个简单的队列。但是,对于希望优先处理或分段处理任务的应用程序来说,将任务推送到多个队列特别有用,因为 Laravel 队列工作进程允许您通过优先级指定应处理哪些队列。例如,如果您将任务推送到 high 队列,则可以运行一个给予其更高处理优先级的工作进程。
1php artisan queue:work --queue=high,default
驱动程序注意事项和先决条件
数据库
为了使用 database 队列驱动程序,您需要一个数据库表来存放任务。通常,这包含在 Laravel 默认的 0001_01_01_000002_create_jobs_table.php 数据库迁移中;但是,如果您的应用程序不包含此迁移,可以使用 make:queue-table Artisan 命令来创建它。
1php artisan make:queue-table2 3php artisan migrate
Redis
为了使用 redis 队列驱动程序,您需要在 config/database.php 配置文件中配置 Redis 数据库连接。
redis 队列驱动程序不支持 serializer 和 compression Redis 选项。
Redis 集群
如果您的 Redis 队列连接使用了 Redis 集群,则您的队列名称必须包含 键哈希标签 (key hash tag)。这是确保给定队列的所有 Redis 键都被放置在同一个哈希槽中所需的。
1'redis' => [2 'driver' => 'redis',3 'connection' => env('REDIS_QUEUE_CONNECTION', 'default'),4 'queue' => env('REDIS_QUEUE', '{default}'),5 'retry_after' => env('REDIS_QUEUE_RETRY_AFTER', 90),6 'block_for' => null,7 'after_commit' => false,8],
阻塞
使用 Redis 队列时,您可以使用 block_for 配置选项来指定驱动程序在循环工作并重新轮询 Redis 数据库之前应等待任务可用的时间。
根据您的队列负载调整此值可能比持续轮询 Redis 数据库获取新任务更有效。例如,您可以将该值设置为 5,表示驱动程序在等待任务可用时应阻塞五秒。
1'redis' => [2 'driver' => 'redis',3 'connection' => env('REDIS_QUEUE_CONNECTION', 'default'),4 'queue' => env('REDIS_QUEUE', 'default'),5 'retry_after' => env('REDIS_QUEUE_RETRY_AFTER', 90),6 'block_for' => 5,7 'after_commit' => false,8],
将 block_for 设置为 0 会导致队列工作进程无限期阻塞,直到任务可用。这也会阻止 SIGTERM 等信号在下个任务处理完成之前得到处理。
其他驱动程序先决条件
所列队列驱动程序需要以下依赖项。可以通过 Composer 包管理器安装这些依赖项。
- Amazon SQS:
aws/aws-sdk-php ~3.0 - Beanstalkd:
pda/pheanstalk ~5.0 - Redis:
predis/predis ~2.0或 phpredis PHP 扩展 - MongoDB:
mongodb/laravel-mongodb
创建任务
生成任务类
默认情况下,应用程序的所有可排队任务都存储在 app/Jobs 目录中。如果 app/Jobs 目录不存在,运行 make:job Artisan 命令时将会创建它。
1php artisan make:job ProcessPodcast
生成的类将实现 Illuminate\Contracts\Queue\ShouldQueue 接口,向 Laravel 表明该任务应被推送到队列中以异步运行。
可以使用 存根发布 (stub publishing) 来自定义任务存根。
类结构
任务类非常简单,通常仅包含一个 handle 方法,该方法在队列处理任务时被调用。首先,让我们看一个任务类的示例。在这个示例中,假设我们管理着一个播客发布服务,需要在播客发布前处理上传的文件。
1<?php 2 3namespace App\Jobs; 4 5use App\Models\Podcast; 6use App\Services\AudioProcessor; 7use Illuminate\Contracts\Queue\ShouldQueue; 8use Illuminate\Foundation\Queue\Queueable; 9 10class ProcessPodcast implements ShouldQueue11{12 use Queueable;13 14 /**15 * Create a new job instance.16 */17 public function __construct(18 public Podcast $podcast,19 ) {}20 21 /**22 * Execute the job.23 */24 public function handle(AudioProcessor $processor): void25 {26 // Process uploaded podcast...27 }28}
在此示例中,请注意我们能够直接将 Eloquent 模型传递到队列任务的构造函数中。由于任务使用了 Queueable trait,当任务处理时,Eloquent 模型及其已加载的关系将优雅地序列化和反序列化。
如果您的队列任务在构造函数中接受 Eloquent 模型,则仅模型的标识符会被序列化到队列中。当任务实际处理时,队列系统会自动从数据库中重新检索完整的模型实例及其已加载的关系。这种模型序列化方法允许向队列驱动程序发送更小的任务负载。
handle 方法依赖注入
handle 方法在队列处理任务时被调用。请注意,我们可以在任务的 handle 方法上进行类型提示依赖项。Laravel 服务容器会自动注入这些依赖项。
如果您想完全控制容器如何将依赖注入到 handle 方法中,可以使用容器的 bindMethod 方法。bindMethod 方法接受一个接收任务和容器的回调。在回调中,您可以随意以任何方式调用 handle 方法。通常,您应该从 App\Providers\AppServiceProvider 服务提供者的 boot 方法中调用此方法。
1use App\Jobs\ProcessPodcast;2use App\Services\AudioProcessor;3use Illuminate\Contracts\Foundation\Application;4 5$this->app->bindMethod([ProcessPodcast::class, 'handle'], function (ProcessPodcast $job, Application $app) {6 return $job->handle($app->make(AudioProcessor::class));7});
二进制数据(例如原始图像内容)在传递给队列任务之前应通过 base64_encode 函数处理。否则,在放入队列时任务可能无法正确序列化为 JSON。
队列关系
由于所有已加载的 Eloquent 模型关系在任务排队时也会被序列化,序列化的任务字符串有时会变得非常大。此外,当任务被反序列化且模型关系从数据库重新检索时,它们将被完整地检索出来。在任务排队过程中序列化模型之前应用的任何先前的关系约束,在任务反序列化时将不会应用。因此,如果您希望处理给定关系的子集,则应该在队列任务中重新约束该关系。
或者,为了防止关系被序列化,您可以在设置属性值时对模型调用 withoutRelations 方法。此方法将返回一个没有加载关系的模型实例。
1/**2 * Create a new job instance.3 */4public function __construct(5 Podcast $podcast,6) {7 $this->podcast = $podcast->withoutRelations();8}
如果您正在使用 PHP 构造函数属性提升 并且希望表明 Eloquent 模型不应序列化其关系,则可以使用 WithoutRelations 属性。
1use Illuminate\Queue\Attributes\WithoutRelations;2 3/**4 * Create a new job instance.5 */6public function __construct(7 #[WithoutRelations]8 public Podcast $podcast,9) {}
为了方便起见,如果您希望序列化所有模型而不包含关系,则可以将 WithoutRelations 属性应用于整个类,而不是将其应用于每个模型。
1<?php 2 3namespace App\Jobs; 4 5use App\Models\DistributionPlatform; 6use App\Models\Podcast; 7use Illuminate\Contracts\Queue\ShouldQueue; 8use Illuminate\Foundation\Queue\Queueable; 9use Illuminate\Queue\Attributes\WithoutRelations;10 11#[WithoutRelations]12class ProcessPodcast implements ShouldQueue13{14 use Queueable;15 16 /**17 * Create a new job instance.18 */19 public function __construct(20 public Podcast $podcast,21 public DistributionPlatform $platform,22 ) {}23}
如果任务接收的是 Eloquent 模型的集合或数组,而不是单个模型,则当任务反序列化并执行时,该集合内的模型将不会恢复它们的关系。这是为了防止处理大量模型的任务占用过多的资源。
唯一任务
唯一任务需要支持 锁 的缓存驱动程序。目前,memcached、redis、dynamodb、database、file 和 array 缓存驱动程序均支持原子锁。
唯一任务约束不适用于批处理中的任务。
有时,您可能希望确保在任何时间点队列中只有一个特定任务的实例。您可以通过在任务类上实现 ShouldBeUnique 接口来做到这一点。此接口不需要您在类上定义任何额外的方法。
1<?php2 3use Illuminate\Contracts\Queue\ShouldQueue;4use Illuminate\Contracts\Queue\ShouldBeUnique;5 6class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique7{8 // ...9}
在上面的示例中,UpdateSearchIndex 任务是唯一的。因此,如果该任务的另一个实例已经在队列中且尚未完成处理,则不会分发该任务。
在某些情况下,您可能希望定义一个使任务唯一的特定“键”,或者您可能希望指定一个超时时间,超过此时间后该任务不再保持唯一。为了实现这一点,您可以使用 UniqueFor 属性并在任务类上定义 uniqueId 方法。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Contracts\Queue\ShouldQueue; 6use Illuminate\Contracts\Queue\ShouldBeUnique; 7use Illuminate\Queue\Attributes\UniqueFor; 8 9#[UniqueFor(3600)]10class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique11{12 /**13 * The product instance.14 *15 * @var \App\Models\Product16 */17 public $product;18 19 /**20 * Get the unique ID for the job.21 */22 public function uniqueId(): string23 {24 return $this->product->id;25 }26}
在上面的示例中,UpdateSearchIndex 任务通过产品 ID 保持唯一。因此,在现有任务完成处理之前,使用相同产品 ID 的任何新任务分发都将被忽略。此外,如果现有任务在一小时内未处理,唯一锁将被释放,具有相同唯一键的另一个任务可以被分发到队列中。
如果您的应用程序从多个 Web 服务器或容器分发任务,则应确保所有服务器都与同一个中央缓存服务器通信,以便 Laravel 能够准确确定任务是否唯一。
保持任务唯一直到处理开始
默认情况下,唯一任务在任务完成处理或尝试完所有重试次数后会“解锁”。但是,在某些情况下,您可能希望任务在处理前立即解锁。为了实现这一点,您的任务应实现 ShouldBeUniqueUntilProcessing 契约,而不是 ShouldBeUnique 契约。
1<?php2 3use Illuminate\Contracts\Queue\ShouldQueue;4use Illuminate\Contracts\Queue\ShouldBeUniqueUntilProcessing;5 6class UpdateSearchIndex implements ShouldQueue, ShouldBeUniqueUntilProcessing7{8 // ...9}
唯一任务锁
在后台,当分发 ShouldBeUnique 任务时,Laravel 会尝试使用 uniqueId 键获取一个 锁。如果锁已被持有,则不会分发任务。当任务完成处理或尝试完所有重试次数后,此锁会被释放。默认情况下,Laravel 将使用默认缓存驱动程序来获取此锁。但是,如果您希望使用另一个驱动程序来获取锁,可以定义一个返回应使用的缓存驱动程序的 uniqueVia 方法。
1use Illuminate\Contracts\Cache\Repository; 2use Illuminate\Support\Facades\Cache; 3 4class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique 5{ 6 // ... 7 8 /** 9 * Get the cache driver for the unique job lock.10 */11 public function uniqueVia(): Repository12 {13 return Cache::driver('redis');14 }15}
如果您只需要限制任务的并发处理,请使用 WithoutOverlapping 任务中间件。
加密任务
Laravel 允许您通过 加密 确保任务数据的隐私和完整性。首先,只需将 ShouldBeEncrypted 接口添加到任务类即可。将此接口添加到类中后,Laravel 会在将任务推送到队列之前自动对其进行加密。
1<?php2 3use Illuminate\Contracts\Queue\ShouldBeEncrypted;4use Illuminate\Contracts\Queue\ShouldQueue;5 6class UpdateSearchIndex implements ShouldQueue, ShouldBeEncrypted7{8 // ...9}
任务中间件
任务中间件允许您在队列任务的执行前后包裹自定义逻辑,从而减少任务本身中的样板代码。例如,考虑以下 handle 方法,它利用 Laravel 的 Redis 速率限制功能来允许每五秒仅处理一个任务:
1use Illuminate\Support\Facades\Redis; 2 3/** 4 * Execute the job. 5 */ 6public function handle(): void 7{ 8 Redis::throttle('key')->block(0)->allow(1)->every(5)->then(function () { 9 info('Lock obtained...');10 11 // Handle job...12 }, function () {13 // Could not obtain lock...14 15 return $this->release(5);16 });17}
虽然此代码有效,但 handle 方法的实现变得嘈杂,因为它被 Redis 速率限制逻辑弄乱了。此外,对于我们想要进行速率限制的任何其他任务,都必须复制此速率限制逻辑。与其在 handle 方法中进行速率限制,我们可以定义一个处理速率限制的任务中间件。
1<?php 2 3namespace App\Jobs\Middleware; 4 5use Closure; 6use Illuminate\Support\Facades\Redis; 7 8class RateLimited 9{10 /**11 * Process the queued job.12 *13 * @param \Closure(object): void $next14 */15 public function handle(object $job, Closure $next): void16 {17 Redis::throttle('key')18 ->block(0)->allow(1)->every(5)19 ->then(function () use ($job, $next) {20 // Lock obtained...21 22 $next($job);23 }, function () use ($job) {24 // Could not obtain lock...25 26 $job->release(5);27 });28 }29}
正如您所见,与 路由中间件 一样,任务中间件会接收正在处理的任务以及应被调用的回调以继续处理该任务。
您可以使用 make:job-middleware Artisan 命令生成一个新的任务中间件类。创建任务中间件后,可以通过从任务的 middleware 方法返回它们来将其附加到任务上。此方法在 make:job Artisan 命令生成的任务中不存在,因此您需要手动将其添加到您的任务类中。
1use App\Jobs\Middleware\RateLimited; 2 3/** 4 * Get the middleware the job should pass through. 5 * 6 * @return array<int, object> 7 */ 8public function middleware(): array 9{10 return [new RateLimited];11}
速率限制
尽管我们刚刚演示了如何编写自己的速率限制任务中间件,但 Laravel 实际上包含了一个速率限制中间件,您可以利用它来限制任务速率。与 路由速率限制器 一样,任务速率限制器使用 RateLimiter 门面的 for 方法定义。
例如,您可能希望允许用户每小时备份一次数据,同时不对高级客户施加此类限制。为了实现这一点,您可以在 AppServiceProvider 的 boot 方法中定义一个 RateLimiter:
1use Illuminate\Cache\RateLimiting\Limit; 2use Illuminate\Support\Facades\RateLimiter; 3 4/** 5 * Bootstrap any application services. 6 */ 7public function boot(): void 8{ 9 RateLimiter::for('backups', function (object $job) {10 return $job->user->vipCustomer()11 ? Limit::none()12 : Limit::perHour(1)->by($job->user->id);13 });14}
在上面的示例中,我们定义了一个每小时的速率限制;但是,您可以使用 perMinute 方法轻松定义基于分钟的速率限制。此外,您可以将任何所需的值传递给速率限制的 by 方法;但是,此值最常用于按客户细分速率限制。
1return Limit::perMinute(50)->by($job->user->id);
一旦定义了速率限制,就可以使用 Illuminate\Queue\Middleware\RateLimited 中间件将速率限制器附加到您的任务。每当任务超过速率限制时,此中间件会将任务释放回队列,并根据速率限制持续时间设置适当的延迟。
1use Illuminate\Queue\Middleware\RateLimited; 2 3/** 4 * Get the middleware the job should pass through. 5 * 6 * @return array<int, object> 7 */ 8public function middleware(): array 9{10 return [new RateLimited('backups')];11}
将速率受限的任务释放回队列仍会增加任务的总 attempts 数。您可能需要相应地调整任务类上的 Tries 和 MaxExceptions 属性。或者,您可能希望使用 retryUntil 方法 来定义该任务不再尝试之前的时间量。
使用 releaseAfter 方法,您还可以指定释放的任务在再次尝试之前必须经过的秒数。
1/**2 * Get the middleware the job should pass through.3 *4 * @return array<int, object>5 */6public function middleware(): array7{8 return [(new RateLimited('backups'))->releaseAfter(60)];9}
如果您不希望在任务受限时重试任务,则可以使用 dontRelease 方法。
1/**2 * Get the middleware the job should pass through.3 *4 * @return array<int, object>5 */6public function middleware(): array7{8 return [(new RateLimited('backups'))->dontRelease()];9}
使用 Redis 进行速率限制
如果您正在使用 Redis,则可以使用 Illuminate\Queue\Middleware\RateLimitedWithRedis 中间件,该中间件针对 Redis 进行了微调,比基本的速率限制中间件效率更高。
1use Illuminate\Queue\Middleware\RateLimitedWithRedis;2 3public function middleware(): array4{5 return [new RateLimitedWithRedis('backups')];6}
connection 方法可用于指定中间件应使用哪个 Redis 连接。
1return [(new RateLimitedWithRedis('backups'))->connection('limiter')];
防止任务重叠
Laravel 包含一个 Illuminate\Queue\Middleware\WithoutOverlapping 中间件,允许您基于任意键防止任务重叠。当队列任务正在修改资源,且该资源同一时间只能由一个任务修改时,这会很有帮助。
例如,假设您有一个更新用户信用评分的队列任务,并且您希望防止同一用户 ID 的信用评分更新任务重叠。为了实现这一点,您可以从任务的 middleware 方法返回 WithoutOverlapping 中间件:
1use Illuminate\Queue\Middleware\WithoutOverlapping; 2 3/** 4 * Get the middleware the job should pass through. 5 * 6 * @return array<int, object> 7 */ 8public function middleware(): array 9{10 return [new WithoutOverlapping($this->user->id)];11}
将重叠任务释放回队列仍会增加任务的总尝试次数。您可能需要相应地调整任务类上的 Tries 和 MaxExceptions 属性。例如,保留默认的 Tries 为 1 将防止任何重叠任务稍后被重试。
相同类型的任何重叠任务都将被释放回队列。您还可以指定释放的任务在再次尝试之前必须经过的秒数。
1/**2 * Get the middleware the job should pass through.3 *4 * @return array<int, object>5 */6public function middleware(): array7{8 return [(new WithoutOverlapping($this->order->id))->releaseAfter(60)];9}
如果您希望立即删除任何重叠的任务,以便不再重试它们,则可以使用 dontRelease 方法。
1/**2 * Get the middleware the job should pass through.3 *4 * @return array<int, object>5 */6public function middleware(): array7{8 return [(new WithoutOverlapping($this->order->id))->dontRelease()];9}
WithoutOverlapping 中间件由 Laravel 的原子锁功能提供支持。有时,您的任务可能会意外失败或超时,导致锁未被释放。因此,您可以使用 expireAfter 方法显式定义锁过期时间。例如,下面的示例将指示 Laravel 在任务开始处理三分钟后释放 WithoutOverlapping 锁:
1/**2 * Get the middleware the job should pass through.3 *4 * @return array<int, object>5 */6public function middleware(): array7{8 return [(new WithoutOverlapping($this->order->id))->expireAfter(180)];9}
WithoutOverlapping 中间件需要支持 锁 的缓存驱动程序。目前,memcached、redis、dynamodb、database、file 和 array 缓存驱动程序均支持原子锁。
跨任务类共享锁键
默认情况下,WithoutOverlapping 中间件仅防止同一类的任务重叠。因此,尽管两个不同的任务类可能使用相同的锁键,但它们不会被阻止重叠。但是,您可以使用 shared 方法指示 Laravel 在任务类之间应用该键:
1use Illuminate\Queue\Middleware\WithoutOverlapping; 2 3class ProviderIsDown 4{ 5 // ... 6 7 public function middleware(): array 8 { 9 return [10 (new WithoutOverlapping("status:{$this->provider}"))->shared(),11 ];12 }13}14 15class ProviderIsUp16{17 // ...18 19 public function middleware(): array20 {21 return [22 (new WithoutOverlapping("status:{$this->provider}"))->shared(),23 ];24 }25}
异常限流
Laravel 包含一个 Illuminate\Queue\Middleware\ThrottlesExceptions 中间件,允许您对异常进行限流。一旦任务抛出指定次数的异常,所有后续执行该任务的尝试都会被延迟,直到指定的时间间隔过去。此中间件对于与不稳定的第三方服务交互的任务特别有用。
例如,假设一个队列任务与开始抛出异常的第三方 API 交互。要对异常进行限流,可以从任务的 middleware 方法返回 ThrottlesExceptions 中间件。通常,此中间件应与实现 基于时间尝试 的任务配对:
1use DateTime; 2use Illuminate\Queue\Middleware\ThrottlesExceptions; 3 4/** 5 * Get the middleware the job should pass through. 6 * 7 * @return array<int, object> 8 */ 9public function middleware(): array10{11 return [new ThrottlesExceptions(10, 5 * 60)];12}13 14/**15 * Determine the time at which the job should timeout.16 */17public function retryUntil(): DateTime18{19 return now()->plus(minutes: 30);20}
中间件接受的第一个构造函数参数是任务被限流前可以抛出的异常数量,第二个构造函数参数是任务被限流后再次尝试之前应经过的秒数。在上面的代码示例中,如果任务抛出 10 个连续异常,我们将等待 5 分钟,然后再尝试该任务,并受 30 分钟的时间限制约束。
当任务抛出异常但尚未达到异常阈值时,通常会立即重试该任务。但是,您可以通过在将中间件附加到任务时调用 backoff 方法来指定此类任务应延迟的分钟数。
1use Illuminate\Queue\Middleware\ThrottlesExceptions; 2 3/** 4 * Get the middleware the job should pass through. 5 * 6 * @return array<int, object> 7 */ 8public function middleware(): array 9{10 return [(new ThrottlesExceptions(10, 5 * 60))->backoff(5)];11}
在内部,此中间件使用 Laravel 的缓存系统来实现速率限制,并且任务类名称被用作缓存“键”。您可以通过在将中间件附加到任务时调用 by 方法来覆盖此键。如果您有多个任务与同一个第三方服务交互,并且希望它们共享一个通用的限流“桶”以确保它们遵守单个共享限制,这将非常有用。
1use Illuminate\Queue\Middleware\ThrottlesExceptions; 2 3/** 4 * Get the middleware the job should pass through. 5 * 6 * @return array<int, object> 7 */ 8public function middleware(): array 9{10 return [(new ThrottlesExceptions(10, 10 * 60))->by('key')];11}
默认情况下,此中间件会限制每个异常。您可以通过在将中间件附加到任务时调用 when 方法来修改此行为。只有当提供给 when 方法的闭包返回 true 时,异常才会被限流。
1use Illuminate\Http\Client\HttpClientException; 2use Illuminate\Queue\Middleware\ThrottlesExceptions; 3 4/** 5 * Get the middleware the job should pass through. 6 * 7 * @return array<int, object> 8 */ 9public function middleware(): array10{11 return [(new ThrottlesExceptions(10, 10 * 60))->when(12 fn (Throwable $throwable) => $throwable instanceof HttpClientException13 )];14}
与 when 方法将任务释放回队列或抛出异常不同,deleteWhen 方法允许您在给定异常发生时完全删除任务。
1use App\Exceptions\CustomerDeletedException; 2use Illuminate\Queue\Middleware\ThrottlesExceptions; 3 4/** 5 * Get the middleware the job should pass through. 6 * 7 * @return array<int, object> 8 */ 9public function middleware(): array10{11 return [(new ThrottlesExceptions(2, 10 * 60))->deleteWhen(CustomerDeletedException::class)];12}
如果您希望将受限异常报告给应用程序的异常处理程序,可以通过在将中间件附加到任务时调用 report 方法来实现。或者,您可以为 report 方法提供一个闭包,只有当给定闭包返回 true 时,才会报告异常。
1use Illuminate\Http\Client\HttpClientException; 2use Illuminate\Queue\Middleware\ThrottlesExceptions; 3 4/** 5 * Get the middleware the job should pass through. 6 * 7 * @return array<int, object> 8 */ 9public function middleware(): array10{11 return [(new ThrottlesExceptions(10, 10 * 60))->report(12 fn (Throwable $throwable) => $throwable instanceof HttpClientException13 )];14}
使用 Redis 进行异常限流
如果您正在使用 Redis,则可以使用 Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis 中间件,该中间件针对 Redis 进行了微调,比基本的异常限流中间件效率更高。
1use Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis;2 3public function middleware(): array4{5 return [new ThrottlesExceptionsWithRedis(10, 10 * 60)];6}
connection 方法可用于指定中间件应使用哪个 Redis 连接。
1return [(new ThrottlesExceptionsWithRedis(10, 10 * 60))->connection('limiter')];
跳过任务
Skip 中间件允许您指定应跳过/删除任务,而无需修改任务的逻辑。如果给定条件评估为 true,Skip::when 方法将删除任务;如果条件评估为 false,Skip::unless 方法将删除任务。
1use Illuminate\Queue\Middleware\Skip; 2 3/** 4 * Get the middleware the job should pass through. 5 */ 6public function middleware(): array 7{ 8 return [ 9 Skip::when($condition),10 ];11}
您也可以将 Closure 传递给 when 和 unless 方法,以进行更复杂的条件评估。
1use Illuminate\Queue\Middleware\Skip; 2 3/** 4 * Get the middleware the job should pass through. 5 */ 6public function middleware(): array 7{ 8 return [ 9 Skip::when(function (): bool {10 return $this->shouldSkip();11 }),12 ];13}
分发任务
编写任务类后,您可以使用任务本身的 dispatch 方法对其进行分发。传递给 dispatch 方法的参数将提供给任务的构造函数。
1<?php 2 3namespace App\Http\Controllers; 4 5use App\Jobs\ProcessPodcast; 6use App\Models\Podcast; 7use Illuminate\Http\RedirectResponse; 8use Illuminate\Http\Request; 9 10class PodcastController extends Controller11{12 /**13 * Store a new podcast.14 */15 public function store(Request $request): RedirectResponse16 {17 $podcast = Podcast::create(/* ... */);18 19 // ...20 21 ProcessPodcast::dispatch($podcast);22 23 return redirect('/podcasts');24 }25}
如果您想有条件地分发任务,可以使用 dispatchIf 和 dispatchUnless 方法。
1ProcessPodcast::dispatchIf($accountActive, $podcast);2 3ProcessPodcast::dispatchUnless($accountSuspended, $podcast);
在新的 Laravel 应用程序中,database 连接被定义为默认队列。您可以通过更改应用程序 .env 文件中的 QUEUE_CONNECTION 环境变量来指定不同的默认队列连接。
延迟分发
如果您想指定任务不应立即供队列工作进程处理,可以在分发任务时使用 delay 方法。例如,让我们指定一个任务在分发 10 分钟后才能进行处理:
1<?php 2 3namespace App\Http\Controllers; 4 5use App\Jobs\ProcessPodcast; 6use App\Models\Podcast; 7use Illuminate\Http\RedirectResponse; 8use Illuminate\Http\Request; 9 10class PodcastController extends Controller11{12 /**13 * Store a new podcast.14 */15 public function store(Request $request): RedirectResponse16 {17 $podcast = Podcast::create(/* ... */);18 19 // ...20 21 ProcessPodcast::dispatch($podcast)22 ->delay(now()->plus(minutes: 10));23 24 return redirect('/podcasts');25 }26}
在某些情况下,任务可能配置了默认延迟。如果您需要绕过此延迟并立即分发任务进行处理,可以使用 withoutDelay 方法。
1ProcessPodcast::dispatch($podcast)->withoutDelay();
Amazon SQS 队列服务的最大延迟时间为 15 分钟。
同步分发
如果您想立即(同步)分发任务,可以使用 dispatchSync 方法。使用此方法时,任务不会被排队,并将在当前进程内立即执行。
1<?php 2 3namespace App\Http\Controllers; 4 5use App\Jobs\ProcessPodcast; 6use App\Models\Podcast; 7use Illuminate\Http\RedirectResponse; 8use Illuminate\Http\Request; 9 10class PodcastController extends Controller11{12 /**13 * Store a new podcast.14 */15 public function store(Request $request): RedirectResponse16 {17 $podcast = Podcast::create(/* ... */);18 19 // Create podcast...20 21 ProcessPodcast::dispatchSync($podcast);22 23 return redirect('/podcasts');24 }25}
延迟分发
使用延迟同步分发,您可以分发一个任务,使其在当前进程中处理,但在 HTTP 响应发送给用户之后。这允许您同步处理“队列”任务,而不会减慢用户的应用程序体验。要延迟同步任务的执行,请将任务分发到 deferred 连接:
1RecordDelivery::dispatch($order)->onConnection('deferred');
deferred 连接也充当默认的 故障转移队列。
同样,background 连接在 HTTP 响应发送给用户后处理任务;但是,该任务是在一个单独派生的 PHP 进程中处理的,允许 PHP-FPM / 应用程序工作进程能够处理另一个传入的 HTTP 请求。
1RecordDelivery::dispatch($order)->onConnection('background');
任务与数据库事务
虽然在数据库事务中分发任务完全没有问题,但您应该特别小心,确保您的任务确实能够成功执行。在事务中分发任务时,任务有可能在父事务提交之前就被工作进程处理。发生这种情况时,您在数据库事务期间对模型或数据库记录所做的任何更新可能尚未在数据库中反映出来。此外,在事务中创建的任何模型或数据库记录在数据库中可能还不存在。
幸运的是,Laravel 提供了几种解决此问题的方法。首先,您可以在队列连接的配置数组中设置 after_commit 连接选项:
1'redis' => [2 'driver' => 'redis',3 // ...4 'after_commit' => true,5],
当 after_commit 选项为 true 时,您可以在数据库事务中分发任务;但是,Laravel 会等到开放的父数据库事务提交后才会实际分发任务。当然,如果当前没有打开任何数据库事务,任务将立即分发。
如果事务由于事务期间发生的异常而回滚,则在该事务期间分发的任务将被丢弃。
将 after_commit 配置选项设置为 true 还会导致任何队列事件监听器、邮件任务、通知和广播事件在所有打开的数据库事务提交后分发。
内联指定提交分发行为
如果您没有将 after_commit 队列连接配置选项设置为 true,您仍然可以指示在所有打开的数据库事务提交后分发特定任务。为了实现这一点,您可以将 afterCommit 方法链接到您的分发操作上:
1use App\Jobs\ProcessPodcast;2 3ProcessPodcast::dispatch($podcast)->afterCommit();
同样,如果 after_commit 配置选项设置为 true,您可以指示立即分发特定任务,而无需等待任何打开的数据库事务提交。
1ProcessPodcast::dispatch($podcast)->beforeCommit();
任务链
任务链允许您指定一系列队列任务,这些任务应在主任务成功执行后按顺序运行。如果序列中的一个任务失败,其余任务将不会运行。要执行队列任务链,可以使用 Bus 门面提供的 chain 方法。Laravel 的命令总线是一个更底层的组件,队列任务分发正是建立在其之上的。
1use App\Jobs\OptimizePodcast; 2use App\Jobs\ProcessPodcast; 3use App\Jobs\ReleasePodcast; 4use Illuminate\Support\Facades\Bus; 5 6Bus::chain([ 7 new ProcessPodcast, 8 new OptimizePodcast, 9 new ReleasePodcast,10])->dispatch();
除了链接任务类实例外,您还可以链接闭包:
1Bus::chain([2 new ProcessPodcast,3 new OptimizePodcast,4 function () {5 Podcast::update(/* ... */);6 },7])->dispatch();
在任务中使用 $this->delete() 方法删除任务不会阻止链接任务的处理。链只有在链中的任务失败时才会停止执行。
链连接与队列
如果您想指定应为链式任务使用的连接和队列,可以使用 onConnection 和 onQueue 方法。除非队列任务被显式分配了不同的连接/队列,否则这些方法将指定应使用的队列连接和队列名称。
1Bus::chain([2 new ProcessPodcast,3 new OptimizePodcast,4 new ReleasePodcast,5])->onConnection('redis')->onQueue('podcasts')->dispatch();
向链中添加任务
有时,您可能需要从链中的另一个任务内向前置或向后追加一个任务到现有的任务链中。您可以使用 prependToChain 和 appendToChain 方法来实现这一点。
1/** 2 * Execute the job. 3 */ 4public function handle(): void 5{ 6 // ... 7 8 // Prepend to the current chain, run job immediately after current job... 9 $this->prependToChain(new TranscribePodcast);10 11 // Append to the current chain, run job at end of chain...12 $this->appendToChain(new TranscribePodcast);13}
链失败
在链接任务时,您可以使用 catch 方法指定如果链中的任务失败时应调用的闭包。给定回调将接收导致任务失败的 Throwable 实例。
1use Illuminate\Support\Facades\Bus; 2use Throwable; 3 4Bus::chain([ 5 new ProcessPodcast, 6 new OptimizePodcast, 7 new ReleasePodcast, 8])->catch(function (Throwable $e) { 9 // A job within the chain has failed...10})->dispatch();
由于链回调是由 Laravel 队列序列化并在稍后执行的,因此不应在链回调中使用 $this 变量。
自定义队列和连接
分发到特定队列
通过将任务推送到不同的队列,您可以对队列任务进行“分类”,甚至可以优先处理分配给不同队列的工作进程数量。请记住,这不会将任务推送到由您的队列配置文件定义的不同的队列“连接”,而只会推送到单个连接内的特定队列。要指定队列,请在分发任务时使用 onQueue 方法。
1<?php 2 3namespace App\Http\Controllers; 4 5use App\Jobs\ProcessPodcast; 6use App\Models\Podcast; 7use Illuminate\Http\RedirectResponse; 8use Illuminate\Http\Request; 9 10class PodcastController extends Controller11{12 /**13 * Store a new podcast.14 */15 public function store(Request $request): RedirectResponse16 {17 $podcast = Podcast::create(/* ... */);18 19 // Create podcast...20 21 ProcessPodcast::dispatch($podcast)->onQueue('processing');22 23 return redirect('/podcasts');24 }25}
或者,您可以通过在任务的构造函数中调用 onQueue 方法来指定任务的队列。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Contracts\Queue\ShouldQueue; 6use Illuminate\Foundation\Queue\Queueable; 7 8class ProcessPodcast implements ShouldQueue 9{10 use Queueable;11 12 /**13 * Create a new job instance.14 */15 public function __construct()16 {17 $this->onQueue('processing');18 }19}
分发到特定连接
如果您的应用程序与多个队列连接交互,您可以使用 onConnection 方法指定将任务推送到哪个连接。
1<?php 2 3namespace App\Http\Controllers; 4 5use App\Jobs\ProcessPodcast; 6use App\Models\Podcast; 7use Illuminate\Http\RedirectResponse; 8use Illuminate\Http\Request; 9 10class PodcastController extends Controller11{12 /**13 * Store a new podcast.14 */15 public function store(Request $request): RedirectResponse16 {17 $podcast = Podcast::create(/* ... */);18 19 // Create podcast...20 21 ProcessPodcast::dispatch($podcast)->onConnection('sqs');22 23 return redirect('/podcasts');24 }25}
您可以将 onConnection 和 onQueue 方法链接在一起,为任务指定连接和队列。
1ProcessPodcast::dispatch($podcast)2 ->onConnection('sqs')3 ->onQueue('processing');
或者,您可以通过在任务的构造函数中调用 onConnection 方法来指定任务的连接。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Contracts\Queue\ShouldQueue; 6use Illuminate\Foundation\Queue\Queueable; 7 8class ProcessPodcast implements ShouldQueue 9{10 use Queueable;11 12 /**13 * Create a new job instance.14 */15 public function __construct()16 {17 $this->onConnection('sqs');18 }19}
队列路由
您可以使用 Queue 门面的 route 方法为特定任务类定义默认连接和队列。当您想要确保某些任务始终使用特定队列,而无需在任务上指定连接或队列时,这非常有用。
除了路由特定任务类外,您还可以将接口、trait 或父类传递给 route 方法。当您这样做时,任何实现该接口、使用该 trait 或扩展该父类的任务都将自动使用配置的连接和队列。
通常,您应该从服务提供者的 boot 方法中调用 route 方法。
1use App\Concerns\RequiresVideo; 2use App\Jobs\ProcessPodcast; 3use App\Jobs\ProcessVideo; 4use Illuminate\Support\Facades\Queue; 5 6/** 7 * Bootstrap any application services. 8 */ 9public function boot(): void10{11 Queue::route(ProcessPodcast::class, connection: 'redis', queue: 'podcasts');12 Queue::route(RequiresVideo::class, queue: 'video');13}
当指定了连接而没有指定队列时,任务将被发送到默认队列。
1Queue::route(ProcessPodcast::class, connection: 'redis');
您还可以通过将数组传递给 route 方法来一次路由多个任务类。
1Queue::route([2 ProcessPodcast::class => ['podcasts', 'redis'], // Queue and connection3 ProcessVideo::class => 'videos', // Queue only (uses default connection)4]);
队列路由仍然可以在每个任务的基础上被任务覆盖。
指定任务最大尝试次数 / 超时值
最大尝试次数
任务尝试是 Laravel 队列系统的核心概念,也是许多高级功能的基础。虽然起初可能会让人感到困惑,但在修改默认配置之前,了解它们的工作方式很重要。
当任务被分发时,它被推送到队列。然后,工作进程会拾取它并尝试执行它。这就是一次任务尝试。
但是,一次尝试并不一定意味着执行了任务的 handle 方法。尝试也可以通过多种方式被“消耗”:
- 任务在执行期间遇到未处理的异常。
- 任务使用
$this->release()手动释放回队列。 - 诸如
WithoutOverlapping或RateLimited之类的中间件未能获取锁并释放了任务。 - 任务超时。
- 任务的
handle方法运行并完成,没有抛出异常。
您可能不希望无限期地尝试任务。因此,Laravel 提供了多种方式来指定任务可以尝试多少次或尝试多长时间。
默认情况下,Laravel 只会尝试一次任务。如果您的任务使用了 WithoutOverlapping 或 RateLimited 等中间件,或者您正在手动释放任务,则很可能需要通过 tries 选项增加允许的尝试次数。
指定任务最大尝试次数的一种方法是通过 Artisan 命令行上的 --tries 开关。这将应用于工作进程处理的所有任务,除非被处理的任务指定了可以尝试的次数。
1php artisan queue:work --tries=3
如果任务超过其最大尝试次数,它将被视为“失败”的任务。有关处理失败任务的更多信息,请参阅 失败任务文档。如果向 queue:work 命令提供了 --tries=0,任务将无限期重试。
您可以采取更细粒度的方法,通过使用 Tries 属性在任务类本身上定义任务可以尝试的最大次数。如果任务上指定了最大尝试次数,它将优先于命令行上提供的 --tries 值。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\Tries; 6 7#[Tries(5)] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
如果您需要动态控制特定任务的最大尝试次数,可以在任务上定义一个 tries 方法。
1/**2 * Determine number of times the job may be attempted.3 */4public function tries(): int5{6 return 5;7}
基于时间的尝试
作为定义任务失败前可尝试次数的替代方案,您可以定义任务不再尝试的时间。这允许任务在给定的时间范围内尝试任意次数。要定义任务不再尝试的时间,请将 retryUntil 方法添加到您的任务类中。此方法应返回一个 DateTime 实例。
1use DateTime;2 3/**4 * Determine the time at which the job should timeout.5 */6public function retryUntil(): DateTime7{8 return now()->plus(minutes: 10);9}
如果同时定义了 retryUntil 和 tries,Laravel 会优先考虑 retryUntil 方法。
最大异常数
有时您可能希望指定任务可以尝试多次,但如果重试是由给定数量的未处理异常触发的(而不是直接由 release 方法释放),则应失败。为了实现这一点,您可以在任务类上使用 Tries 和 MaxExceptions 属性。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Contracts\Queue\ShouldQueue; 6use Illuminate\Foundation\Queue\Queueable; 7use Illuminate\Queue\Attributes\MaxExceptions; 8use Illuminate\Queue\Attributes\Tries; 9use Illuminate\Support\Facades\Redis;10 11#[Tries(25)]12#[MaxExceptions(3)]13class ProcessPodcast implements ShouldQueue14{15 use Queueable;16 17 /**18 * Execute the job.19 */20 public function handle(): void21 {22 Redis::throttle('key')->allow(10)->every(60)->then(function () {23 // Lock obtained, process the podcast...24 }, function () {25 // Unable to obtain lock...26 return $this->release(10);27 });28 }29}
在此示例中,如果应用程序无法获取 Redis 锁,任务将释放十秒,并将继续重试最多 25 次。但是,如果任务抛出三个未处理的异常,则任务将失败。
超时
通常,您大致知道队列任务需要花费多长时间。因此,Laravel 允许您指定“超时”值。默认情况下,超时值为 60 秒。如果任务处理时间超过超时值指定的秒数,处理任务的工作进程将退出并报错。通常,工作进程将由 在服务器上配置的进程管理器 自动重启。
可以使用 Artisan 命令行上的 --timeout 开关指定任务可以运行的最大秒数。
1php artisan queue:work --timeout=30
如果任务因持续超时而超过其最大尝试次数,它将被标记为失败。
您还可以使用任务类上的 Timeout 属性定义允许任务运行的最大秒数。如果任务上指定了超时,它将优先于命令行上指定的任何超时。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\Timeout; 6 7#[Timeout(120)] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
有时,套接字或传出 HTTP 连接等 IO 阻塞进程可能不遵守您指定的超时。因此,在使用这些功能时,您也应始终尝试使用它们的 API 指定超时。例如,使用 Guzzle 时,应始终指定连接和请求超时值。
超时失败
如果您想指示任务在超时时应被标记为 失败,可以使用任务类上的 FailOnTimeout 属性:
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\FailOnTimeout; 6 7#[FailOnTimeout] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
默认情况下,当任务超时时,它会消耗一次尝试并被释放回队列(如果允许重试)。但是,如果您将任务配置为超时失败,则无论为 tries 设置的值如何,它都不会被重试。
SQS FIFO 和公平队列
Laravel 支持 Amazon SQS FIFO (先进先出) 队列,允许您以发送的确切顺序处理任务,同时通过消息去重确保有且仅有一次的处理。
FIFO 队列需要一个消息组 ID 来确定哪些任务可以并行处理。具有相同组 ID 的任务按顺序处理,而具有不同组 ID 的消息可以并发处理。
Laravel 提供了一个流式 onGroup 方法,用于在分发任务时指定消息组 ID。
1ProcessOrder::dispatch($order)2 ->onGroup("customer-{$order->customer_id}");
SQS FIFO 队列支持消息去重以确保有且仅有一次的处理。在您的任务类中实现一个 deduplicationId 方法来提供自定义去重 ID。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Contracts\Queue\ShouldQueue; 6use Illuminate\Foundation\Queue\Queueable; 7 8class ProcessSubscriptionRenewal implements ShouldQueue 9{10 use Queueable;11 12 // ...13 14 /**15 * Get the job's deduplication ID.16 */17 public function deduplicationId(): string18 {19 return "renewal-{$this->subscription->id}";20 }21}
FIFO 监听器、邮件和通知
使用 FIFO 队列时,您还需要在监听器、邮件和通知上定义消息组。或者,您可以将这些对象的可队列实例分发到非 FIFO 队列。
要为 队列事件监听器 定义消息组,请在监听器上定义一个 messageGroup 方法。您也可以选择定义一个 deduplicationId 方法。
1<?php 2 3namespace App\Listeners; 4 5class SendShipmentNotification 6{ 7 // ... 8 9 /**10 * Get the job's message group.11 */12 public function messageGroup(): string13 {14 return 'shipments';15 }16 17 /**18 * Get the job's deduplication ID.19 */20 public function deduplicationId(): string21 {22 return "shipment-notification-{$this->shipment->id}";23 }24}
发送将在 FIFO 队列上排队的 邮件消息 时,您应该在发送通知时调用 onGroup 方法,并可选择调用 withDeduplicator 方法。
1use App\Mail\InvoicePaid;2use Illuminate\Support\Facades\Mail;3 4$invoicePaid = (new InvoicePaid($invoice))5 ->onGroup('invoices')6 ->withDeduplicator(fn () => 'invoices-'.$invoice->id);7 8Mail::to($request->user())->send($invoicePaid);
发送将在 FIFO 队列上排队的 通知 时,您应该在发送通知时调用 onGroup 方法,并可选择调用 withDeduplicator 方法。
1use App\Notifications\InvoicePaid;2 3$invoicePaid = (new InvoicePaid($invoice))4 ->onGroup('invoices')5 ->withDeduplicator(fn () => 'invoices-'.$invoice->id);6 7$user->notify($invoicePaid);
队列故障转移
failover 队列驱动程序在将任务推送到队列时提供自动故障转移功能。如果 failover 配置的主队列连接因任何原因失败,Laravel 将自动尝试将任务推送到列表中配置的下一个连接。这对于确保生产环境中的高可用性特别有用,其中队列可靠性至关重要。
要配置故障转移队列连接,请指定 failover 驱动程序并提供一个连接名称数组,以按顺序尝试。默认情况下,Laravel 在应用程序的 config/queue.php 配置文件中包含了一个故障转移配置示例。
1'failover' => [2 'driver' => 'failover',3 'connections' => [4 'redis',5 'database',6 'sync',7 ],8],
配置了使用 failover 驱动程序的连接后,您需要在应用程序的 .env 文件中将故障转移连接设置为默认队列连接,以使用故障转移功能。
1QUEUE_CONNECTION=failover
接下来,为故障转移连接列表中的每个连接启动至少一个工作进程。
1php artisan queue:work redis2php artisan queue:work database
您不需要为使用 sync、background 或 deferred 队列驱动程序的连接运行工作进程,因为这些驱动程序在当前 PHP 进程中处理任务。
当队列连接操作失败且激活故障转移时,Laravel 将分发 Illuminate\Queue\Events\QueueFailedOver 事件,允许您报告或记录队列连接已失败。
如果您使用 Laravel Horizon,请记住 Horizon 仅管理 Redis 队列。如果您的故障转移列表包含 database,您应该与 Horizon 一起运行常规的 php artisan queue:work database 进程。
错误处理
如果任务处理期间抛出异常,该任务将自动释放回队列,以便再次尝试。任务将继续释放,直到尝试次数达到应用程序允许的最大次数。最大尝试次数由 queue:work Artisan 命令上使用的 --tries 开关定义。或者,可以在任务类本身上定义最大尝试次数。有关运行队列工作进程的更多信息 可以在下方找到。
手动释放任务
有时您可能希望手动将任务释放回队列,以便稍后再次尝试。您可以通过调用 release 方法来实现这一点:
1/**2 * Execute the job.3 */4public function handle(): void5{6 // ...7 8 $this->release();9}
默认情况下,release 方法会将任务释放回队列以进行立即处理。但是,您可以通过将整数或日期实例传递给 release 方法,来指示队列在给定的秒数过去之前不使任务可供处理。
1$this->release(10);2 3$this->release(now()->plus(seconds: 10));
手动使任务失败
有时您可能需要手动将任务标记为“失败”。为此,您可以调用 fail 方法:
1/**2 * Execute the job.3 */4public function handle(): void5{6 // ...7 8 $this->fail();9}
如果您想因捕获到的异常而将任务标记为失败,可以将该异常传递给 fail 方法。或者,为了方便起见,您可以传递一个字符串错误消息,该消息将为您转换为异常。
1$this->fail($exception);2 3$this->fail('Something went wrong.');
有关失败任务的更多信息,请查看 处理任务失败的文档。
在特定异常上使任务失败
FailOnException 任务中间件 允许您在抛出特定异常时短路重试。这允许对瞬态异常(如外部 API 错误)进行重试,但对于持久异常(如用户权限被撤销)使任务永久失败。
1<?php 2 3namespace App\Jobs; 4 5use App\Models\User; 6use Illuminate\Auth\Access\AuthorizationException; 7use Illuminate\Contracts\Queue\ShouldQueue; 8use Illuminate\Foundation\Queue\Queueable; 9use Illuminate\Queue\Attributes\Tries;10use Illuminate\Queue\Middleware\FailOnException;11use Illuminate\Support\Facades\Http;12 13#[Tries(3)]14class SyncChatHistory implements ShouldQueue15{16 use Queueable;17 18 /**19 * Create a new job instance.20 */21 public function __construct(22 public User $user,23 ) {}24 25 /**26 * Execute the job.27 */28 public function handle(): void29 {30 $this->user->authorize('sync-chat-history');31 32 $response = Http::throw()->get(33 "https://chat.laravel.test/?user={$this->user->uuid}"34 );35 36 // ...37 }38 39 /**40 * Get the middleware the job should pass through.41 */42 public function middleware(): array43 {44 return [45 new FailOnException([AuthorizationException::class])46 ];47 }48}
任务批处理
Laravel 的任务批处理功能允许您轻松执行一组并行任务,然后在批处理任务执行完成后执行某些操作。
在开始之前,您应该创建一个数据库迁移来构建一个表,其中将包含有关任务批处理的元信息,例如它们的完成百分比。此迁移可以使用 make:queue-batches-table Artisan 命令生成。
1php artisan make:queue-batches-table2 3php artisan migrate
定义可批处理的任务
要定义可批处理的任务,您应该像往常一样 创建一个可队列的任务;但是,您应该将 Illuminate\Bus\Batchable trait 添加到任务类中。此 trait 提供了对 batch 方法的访问,可用于检索任务在其内执行的当前批处理。
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Bus\Batchable; 6use Illuminate\Contracts\Queue\ShouldQueue; 7use Illuminate\Foundation\Queue\Queueable; 8 9class ImportCsv implements ShouldQueue10{11 use Batchable, Queueable;12 13 /**14 * Execute the job.15 */16 public function handle(): void17 {18 if ($this->batch()->cancelled()) {19 // Determine if the batch has been cancelled...20 21 return;22 }23 24 // Import a portion of the CSV file...25 }26}
分发批处理
要分发任务批处理,您应该使用 Bus 门面的 batch 方法。当然,批处理主要在与完成回调结合使用时才有用。因此,您可以使用 then、catch 和 finally 方法为批处理定义完成回调。当调用这些回调时,它们中的每一个都将接收一个 Illuminate\Bus\Batch 实例。
运行多个队列工作进程时,批处理中的任务将并行处理。因此,任务完成的顺序可能与添加到批处理的顺序不一致。有关如何按顺序运行一系列任务的信息,请咨询我们关于 任务链和批处理 的文档。
在此示例中,我们将想象我们正在排队处理一批任务,每个任务处理 CSV 文件中的给定行数:
1use App\Jobs\ImportCsv; 2use Illuminate\Bus\Batch; 3use Illuminate\Support\Facades\Bus; 4use Throwable; 5 6$batch = Bus::batch([ 7 new ImportCsv(1, 100), 8 new ImportCsv(101, 200), 9 new ImportCsv(201, 300),10 new ImportCsv(301, 400),11 new ImportCsv(401, 500),12])->before(function (Batch $batch) {13 // The batch has been created but no jobs have been added...14})->progress(function (Batch $batch) {15 // A single job has completed successfully...16})->then(function (Batch $batch) {17 // All jobs completed successfully...18})->catch(function (Batch $batch, Throwable $e) {19 // Batch job failure detected...20})->finally(function (Batch $batch) {21 // The batch has finished executing...22})->dispatch();23 24return $batch->id;
可以通过 $batch->id 属性访问的批处理 ID,可用于在批处理分发后 查询 Laravel 命令总线 以获取有关批处理的信息。
由于批处理回调是由 Laravel 队列序列化并在稍后执行的,因此不应在回调中使用 $this 变量。此外,由于批处理任务被包裹在数据库事务中,触发隐式提交的数据库语句不应在任务中执行。
命名批处理
如果命名了批处理,某些工具(如 Laravel Horizon 和 Laravel Telescope)可能会为批处理提供更用户友好的调试信息。要为批处理分配任意名称,您可以在定义批处理时调用 name 方法。
1$batch = Bus::batch([2 // ...3])->then(function (Batch $batch) {4 // All jobs completed successfully...5})->name('Import CSV')->dispatch();
批处理连接与队列
如果您想指定应为批处理任务使用的连接和队列,可以使用 onConnection 和 onQueue 方法。所有批处理任务必须在同一个连接和队列内执行。
1$batch = Bus::batch([2 // ...3])->then(function (Batch $batch) {4 // All jobs completed successfully...5})->onConnection('redis')->onQueue('imports')->dispatch();
链与批处理
您可以通过将链式任务放入数组中,在批处理中定义一组 链式任务。例如,我们可以并行执行两个任务链,并在两个任务链都完成处理后执行回调:
1use App\Jobs\ReleasePodcast; 2use App\Jobs\SendPodcastReleaseNotification; 3use Illuminate\Bus\Batch; 4use Illuminate\Support\Facades\Bus; 5 6Bus::batch([ 7 [ 8 new ReleasePodcast(1), 9 new SendPodcastReleaseNotification(1),10 ],11 [12 new ReleasePodcast(2),13 new SendPodcastReleaseNotification(2),14 ],15])->then(function (Batch $batch) {16 // All jobs completed successfully...17})->dispatch();
相反,您可以通过在链中定义批处理来在 链 中运行任务批处理。例如,您可以先运行一批任务来发布多个播客,然后运行一批任务来发送发布通知:
1use App\Jobs\FlushPodcastCache; 2use App\Jobs\ReleasePodcast; 3use App\Jobs\SendPodcastReleaseNotification; 4use Illuminate\Support\Facades\Bus; 5 6Bus::chain([ 7 new FlushPodcastCache, 8 Bus::batch([ 9 new ReleasePodcast(1),10 new ReleasePodcast(2),11 ]),12 Bus::batch([13 new SendPodcastReleaseNotification(1),14 new SendPodcastReleaseNotification(2),15 ]),16])->dispatch();
向批处理添加任务
有时从批处理任务内向批处理添加额外任务可能会很有用。当您需要批处理成千上万个任务,而这些任务在 Web 请求期间分发花费太长时间时,这种模式很有用。因此,您可以分发初始的“加载程序”任务批处理,这些任务将用更多任务填充批处理。
1$batch = Bus::batch([2 new LoadImportBatch,3 new LoadImportBatch,4 new LoadImportBatch,5])->then(function (Batch $batch) {6 // All jobs completed successfully...7})->name('Import Contacts')->dispatch();
在此示例中,我们将使用 LoadImportBatch 任务来用额外任务填充批处理。为了实现这一点,我们可以使用可以通过任务的 batch 方法访问的批处理实例上的 add 方法:
1use App\Jobs\ImportContacts; 2use Illuminate\Support\Collection; 3 4/** 5 * Execute the job. 6 */ 7public function handle(): void 8{ 9 if ($this->batch()->cancelled()) {10 return;11 }12 13 $this->batch()->add(Collection::times(1000, function () {14 return new ImportContacts;15 }));16}
您只能从属于同一批处理的任务内向批处理添加任务。
检查批处理
提供给批处理完成回调的 Illuminate\Bus\Batch 实例具有各种属性和方法,可帮助您与给定的任务批处理交互并检查它们。
1// The UUID of the batch... 2$batch->id; 3 4// The name of the batch (if applicable)... 5$batch->name; 6 7// The number of jobs assigned to the batch... 8$batch->totalJobs; 9 10// The number of jobs that have not been processed by the queue...11$batch->pendingJobs;12 13// The number of jobs that have failed...14$batch->failedJobs;15 16// The number of jobs that have been processed thus far...17$batch->processedJobs();18 19// The completion percentage of the batch (0-100)...20$batch->progress();21 22// Indicates if the batch has finished executing...23$batch->finished();24 25// Cancel the execution of the batch...26$batch->cancel();27 28// Indicates if the batch has been cancelled...29$batch->cancelled();
从路由返回批处理
所有 Illuminate\Bus\Batch 实例都是 JSON 可序列化的,这意味着您可以直接从应用程序的路由之一返回它们,以检索包含有关批处理信息(包括其完成进度)的 JSON 有效负载。这使得在应用程序 UI 中显示有关批处理完成进度的信息变得很方便。
要按 ID 检索批处理,可以使用 Bus 门面的 findBatch 方法:
1use Illuminate\Support\Facades\Bus;2use Illuminate\Support\Facades\Route;3 4Route::get('/batch/{batchId}', function (string $batchId) {5 return Bus::findBatch($batchId);6});
取消批处理
有时您可能需要取消给定批处理的执行。可以通过在 Illuminate\Bus\Batch 实例上调用 cancel 方法来实现这一点:
1/** 2 * Execute the job. 3 */ 4public function handle(): void 5{ 6 if ($this->user->exceedsImportLimit()) { 7 $this->batch()->cancel(); 8 9 return;10 }11 12 if ($this->batch()->cancelled()) {13 return;14 }15}
正如您在前面的示例中可能注意到的那样,批处理任务通常应在继续执行之前确定其相应的批处理是否已取消。但是,为了方便起见,您可以改为将 SkipIfBatchCancelled 中间件 分配给任务。顾名思义,此中间件将指示 Laravel 在其相应的批处理已取消的情况下不处理该任务。
1use Illuminate\Queue\Middleware\SkipIfBatchCancelled;2 3/**4 * Get the middleware the job should pass through.5 */6public function middleware(): array7{8 return [new SkipIfBatchCancelled];9}
批处理失败
当批处理任务失败时,将调用 catch 回调(如果已分配)。此回调仅针对批处理中第一个失败的任务调用。
允许失败
当批处理中的任务失败时,Laravel 会自动将批处理标记为“已取消”。如果您愿意,可以禁用此行为,以便任务失败不会自动将批处理标记为已取消。这可以通过在分发批处理时调用 allowFailures 方法来实现:
1$batch = Bus::batch([2 // ...3])->then(function (Batch $batch) {4 // All jobs completed successfully...5})->allowFailures()->dispatch();
您可以选择为 allowFailures 方法提供一个闭包,该闭包将在每次任务失败时执行。
1$batch = Bus::batch([2 // ...3])->allowFailures(function (Batch $batch, $exception) {4 // Handle individual job failures...5})->dispatch();
重试失败的批处理任务
为了方便起见,Laravel 提供了一个 queue:retry-batch Artisan 命令,允许您轻松重试给定批处理的所有失败任务。此命令接受应重试其失败任务的批处理的 UUID。
1php artisan queue:retry-batch 32dbc76c-4f82-4749-b610-a639fe0099b5
清理批处理
如果不清理,job_batches 表可以非常快地积累记录。为了缓解这种情况,您应该 调度 queue:prune-batches Artisan 命令每天运行。
1use Illuminate\Support\Facades\Schedule;2 3Schedule::command('queue:prune-batches')->daily();
默认情况下,所有超过 24 小时的已完成批处理都将被清理。您可以在调用命令时使用 hours 选项来确定保留批处理数据的时间。例如,以下命令将删除所有在 48 小时前完成的批处理:
1use Illuminate\Support\Facades\Schedule;2 3Schedule::command('queue:prune-batches --hours=48')->daily();
有时,您的 job_batches 表可能会积累从未成功完成的批处理记录,例如任务失败且该任务从未成功重试的批处理。您可以指示 queue:prune-batches 命令使用 unfinished 选项来清理这些未完成的批处理记录。
1use Illuminate\Support\Facades\Schedule;2 3Schedule::command('queue:prune-batches --hours=48 --unfinished=72')->daily();
同样,您的 job_batches 表也可能积累已取消批处理的批处理记录。您可以指示 queue:prune-batches 命令使用 cancelled 选项来清理这些已取消的批处理记录。
1use Illuminate\Support\Facades\Schedule;2 3Schedule::command('queue:prune-batches --hours=48 --cancelled=72')->daily();
在 DynamoDB 中存储批处理
Laravel 还提供支持将批处理元信息存储在 DynamoDB 中,而不是关系数据库中。但是,您需要手动创建一个 DynamoDB 表来存储所有批处理记录。
通常,此表应命名为 job_batches,但您应该根据应用程序 queue 配置文件中的 queue.batching.table 配置值来命名该表。
DynamoDB 批处理表配置
job_batches 表应具有一个名为 application 的字符串主分区键和一个名为 id 的字符串主排序键。键的 application 部分将包含您的应用程序名称,该名称由应用程序 app 配置文件中的 name 配置值定义。由于应用程序名称是 DynamoDB 表键的一部分,因此您可以使用同一个表来存储多个 Laravel 应用程序的任务批处理。
此外,如果您想利用 自动批处理清理,可以为您的表定义 ttl 属性。
DynamoDB 配置
接下来,安装 AWS SDK,以便您的 Laravel 应用程序可以与 Amazon DynamoDB 通信:
1composer require aws/aws-sdk-php
然后,将 queue.batching.driver 配置选项的值设置为 dynamodb。此外,您应该在 batching 配置数组中定义 key、secret 和 region 配置选项。这些选项将用于向 AWS 进行身份验证。使用 dynamodb 驱动程序时,queue.batching.database 配置选项是不必要的。
1'batching' => [2 'driver' => env('QUEUE_BATCHING_DRIVER', 'dynamodb'),3 'key' => env('AWS_ACCESS_KEY_ID'),4 'secret' => env('AWS_SECRET_ACCESS_KEY'),5 'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),6 'table' => 'job_batches',7],
在 DynamoDB 中清理批处理
利用 DynamoDB 存储任务批处理信息时,用于清理存储在关系数据库中的批处理的常规清理命令将不起作用。相反,您可以利用 DynamoDB 的原生 TTL 功能 来自动删除旧批处理的记录。
如果您使用 ttl 属性定义了 DynamoDB 表,则可以定义配置参数来指示 Laravel 如何清理批处理记录。queue.batching.ttl_attribute 配置值定义了保存 TTL 的属性名称,而 queue.batching.ttl 配置值定义了批处理记录相对于记录最后更新时间,在多少秒后可以从 DynamoDB 表中删除。
1'batching' => [2 'driver' => env('QUEUE_FAILED_DRIVER', 'dynamodb'),3 'key' => env('AWS_ACCESS_KEY_ID'),4 'secret' => env('AWS_SECRET_ACCESS_KEY'),5 'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),6 'table' => 'job_batches',7 'ttl_attribute' => 'ttl',8 'ttl' => 60 * 60 * 24 * 7, // 7 days...9],
队列闭包
除了将任务类分发到队列外,您还可以分发一个闭包。这非常适合需要在当前请求周期之外执行的快速、简单的任务。当将闭包分发到队列时,闭包的代码内容会被加密签名,以确保其在传输过程中不会被修改。
1use App\Models\Podcast;2 3$podcast = Podcast::find(1);4 5dispatch(function () use ($podcast) {6 $podcast->publish();7});
要为排队的闭包分配名称(队列报告仪表盘可以使用该名称,并且 queue:work 命令也会显示该名称),您可以使用 name 方法:
1dispatch(function () {2 // ...3})->name('Publish Podcast');
使用 catch 方法,您可以提供一个闭包,如果在用尽所有队列 配置的重试尝试 后,排队的闭包未能成功完成,则应执行该闭包:
1use Throwable;2 3dispatch(function () use ($podcast) {4 $podcast->publish();5})->catch(function (Throwable $e) {6 // This job has failed...7});
由于 catch 回调是由 Laravel 队列序列化并在稍后执行的,因此不应在 catch 回调中使用 $this 变量。
运行队列工作进程
queue:work 命令
Laravel 包含一个 Artisan 命令,它将启动队列工作进程并在新任务被推送到队列时处理它们。您可以使用 queue:work Artisan 命令运行该工作进程。请注意,一旦 queue:work 命令启动,它将持续运行,直到被手动停止或您关闭终端。
1php artisan queue:work
要使 queue:work 进程在后台永久运行,您应该使用进程监控器(例如 Supervisor)来确保队列工作进程不会停止运行。
如果您希望在命令输出中包含已处理的任务 ID、连接名称和队列名称,可以在调用 queue:work 命令时包含 -v 标志。
1php artisan queue:work -v
请记住,队列工作进程是长期运行的进程,并将启动的应用程序状态存储在内存中。因此,它们在启动后不会注意到代码库中的更改。因此,在您的部署过程中,请务必 重启您的队列工作进程。此外,请记住应用程序创建或修改的任何静态状态都不会在任务之间自动重置。
或者,您可以运行 queue:listen 命令。使用 queue:listen 命令时,无需在想要重新加载更新的代码或重置应用程序状态时手动重启工作进程;但是,此命令的效率明显低于 queue:work 命令。
1php artisan queue:listen
运行多个队列工作进程
要向队列分配多个工作进程并并发处理任务,您只需启动多个 queue:work 进程。这可以在本地通过终端中的多个选项卡完成,或者在生产环境中使用进程管理器的配置设置完成。使用 Supervisor 时,可以使用 numprocs 配置值。
指定连接和队列
您还可以指定工作进程应利用哪个队列连接。传递给 work 命令的连接名称应对应于 config/queue.php 配置文件中定义的连接之一。
1php artisan queue:work redis
默认情况下,queue:work 命令仅处理给定连接上默认队列的任务。但是,您可以通过仅处理给定连接的特定队列来进一步自定义队列工作进程。例如,如果您的所有电子邮件都在 redis 队列连接的 emails 队列中处理,则可以发出以下命令来启动仅处理该队列的工作进程:
1php artisan queue:work redis --queue=emails
处理指定数量的任务
--once 选项可用于指示工作进程仅处理队列中的单个任务:
1php artisan queue:work --once
--max-jobs 选项可用于指示工作进程处理给定数量的任务然后退出。当与 Supervisor 结合使用时,此选项可能很有用,这样您的工作进程在处理给定数量的任务后会自动重启,从而释放它们可能积累的任何内存。
1php artisan queue:work --max-jobs=1000
处理所有排队任务然后退出
--stop-when-empty 选项可用于指示工作进程处理所有任务然后优雅地退出。当您希望在队列为空后关闭容器时,此选项在 Docker 容器内处理 Laravel 队列时非常有用。
1php artisan queue:work --stop-when-empty
处理给定秒数的任务
--max-time 选项可用于指示工作进程处理任务给定的秒数,然后退出。当与 Supervisor 结合使用时,此选项可能很有用,这样您的工作进程在处理给定时间的任务后会自动重启,从而释放它们可能积累的任何内存。
1# Process jobs for one hour and then exit...2php artisan queue:work --max-time=3600
工作进程休眠时间
当队列中有任务可用时,工作进程将持续处理任务,任务之间没有延迟。但是,sleep 选项决定了如果没有任务可用,工作进程将“休眠”多少秒。当然,在休眠时,工作进程不会处理任何新任务。
1php artisan queue:work --sleep=3
维护模式与队列
当您的应用程序处于 维护模式 时,不会处理任何队列任务。一旦应用程序退出维护模式,任务将照常处理。
要强制您的队列工作进程即使在启用了维护模式的情况下也处理任务,可以使用 --force 选项:
1php artisan queue:work --force
资源注意事项
守护进程队列工作程序(Worker)在处理每个作业(Job)之前不会“重启”框架。因此,在每个作业完成后,你应该释放任何繁重的资源。例如,如果你正在使用 GD 库进行图像处理,那么在处理完图像后,你应该使用 imagedestroy 释放内存。
队列优先级
有时你可能希望确定队列的处理优先级。例如,在 config/queue.php 配置文件中,你可以将 redis 连接的默认 queue 设置为 low。但是,有时你可能希望将作业推送到 high 优先级的队列中,如下所示:
1dispatch((new Job)->onQueue('high'));
要启动一个工作程序,以确保在继续处理 low 队列中的任何作业之前先处理完 high 队列中的所有作业,请将以逗号分隔的队列名称列表传递给 work 命令:
1php artisan queue:work --queue=high,low
队列工作进程与部署
由于队列工作程序是长生命周期的进程,如果不重启它们,它们将无法感知代码的变更。因此,部署使用队列工作程序的应用程序最简单的方法是在部署过程中重启工作程序。你可以通过执行 queue:restart 命令来优雅地重启所有工作程序:
1php artisan queue:restart
该命令将指示所有队列工作程序在处理完当前作业后优雅地退出,从而确保不会丢失现有的作业。由于执行 queue:restart 命令时队列工作程序会退出,因此你应该运行诸如 Supervisor 之类的进程管理器来自动重启队列工作程序。
队列使用 缓存 来存储重启信号,因此在使用此功能之前,请确保已为应用程序正确配置了缓存驱动程序。
任务过期与超时
作业过期
在 config/queue.php 配置文件中,每个队列连接都定义了一个 retry_after 选项。此选项指定了队列连接在重试正在处理的作业之前应等待的秒数。例如,如果 retry_after 的值设置为 90,那么如果作业处理了 90 秒而没有被释放或删除,它将被重新放回队列。通常,你应该将 retry_after 的值设置为你的作业完成处理所需的最长秒数。
唯一不包含 retry_after 值的队列连接是 Amazon SQS。SQS 将根据在 AWS 控制台中管理的 默认可见性超时 (Default Visibility Timeout) 来重试作业。
工作程序超时
queue:work Artisan 命令提供了一个 --timeout 选项。默认情况下,--timeout 的值为 60 秒。如果作业的处理时间超过了超时值指定的秒数,处理该作业的工作程序将退出并报错。通常,工作程序会由服务器上配置的进程管理器自动重启。
1php artisan queue:work --timeout=60
retry_after 配置选项和 --timeout CLI 选项是不同的,但它们协同工作以确保作业不会丢失,并且每个作业仅被成功处理一次。
--timeout 的值应始终至少比你的 retry_after 配置值短几秒。这将确保处理卡死作业的工作程序在作业被重试之前终止。如果你的 --timeout 选项比 retry_after 配置值长,你的作业可能会被处理两次。
暂停和恢复队列工作进程
有时你可能需要临时阻止队列工作程序处理新作业,而不必完全停止工作程序。例如,你可能希望在系统维护期间暂停作业处理。Laravel 提供了 queue:pause 和 queue:continue Artisan 命令来暂停和恢复队列工作程序。
要暂停特定的队列,请提供队列连接名称和队列名称:
1php artisan queue:pause database:default
在此示例中,database 是队列连接名称,default 是队列名称。一旦队列被暂停,处理来自该队列作业的任何工作程序将继续完成其当前作业,但在队列恢复之前不会获取任何新作业。
要恢复处理暂停队列上的作业,请使用 queue:continue 命令:
1php artisan queue:continue database:default
恢复队列后,工作程序将立即开始处理来自该队列的新作业。请注意,暂停队列不会停止工作程序进程本身——它只会阻止工作程序从指定队列处理新作业。
工作程序重启和暂停信号
默认情况下,队列工作程序会在每次作业迭代时轮询缓存驱动程序以获取重启和暂停信号。虽然这种轮询对于响应 queue:restart 和 queue:pause 命令至关重要,但它确实会带来少量的性能开销。
如果你需要优化性能且不需要这些中断功能,可以通过在 Queue 门面(Facade)上调用 withoutInterruptionPolling 方法来全局禁用此轮询。这通常应该在 AppServiceProvider 的 boot 方法中完成:
1use Illuminate\Support\Facades\Queue;2 3/**4 * Bootstrap any application services.5 */6public function boot(): void7{8 Queue::withoutInterruptionPolling();9}
或者,你可以通过在 Illuminate\Queue\Worker 类上设置静态的 $restartable 或 $pausable 属性,分别禁用重启或暂停轮询:
1use Illuminate\Queue\Worker; 2 3/** 4 * Bootstrap any application services. 5 */ 6public function boot(): void 7{ 8 Worker::$restartable = false; 9 Worker::$pausable = false;10}
当禁用中断轮询时,工作程序将不会响应 queue:restart 或 queue:pause 命令(取决于禁用了哪些功能)。
Supervisor 配置
在生产环境中,你需要一种方法来保持 queue:work 进程运行。queue:work 进程可能会因各种原因停止运行,例如工作程序超时或执行了 queue:restart 命令。
因此,你需要配置一个进程监视器,它可以在 queue:work 进程退出时检测到并自动重启它们。此外,进程监视器允许你指定希望同时运行多少个 queue:work 进程。Supervisor 是 Linux 环境中常用的进程监视器,我们将在接下来的文档中讨论如何配置它。
安装 Supervisor
Supervisor 是 Linux 操作系统的一个进程监视器,如果你的 queue:work 进程失败,它会自动重启它们。要在 Ubuntu 上安装 Supervisor,你可以使用以下命令:
1sudo apt-get install supervisor
如果自行配置和管理 Supervisor 对你来说太复杂,请考虑使用 Laravel Cloud,它提供了一个完全托管的平台来运行 Laravel 队列工作程序。
配置 Supervisor
Supervisor 配置文件通常存储在 /etc/supervisor/conf.d 目录中。在此目录中,你可以创建任意数量的配置文件,以指示 Supervisor 如何监视你的进程。例如,让我们创建一个 laravel-worker.conf 文件来启动和监视 queue:work 进程:
1[program:laravel-worker] 2process_name=%(program_name)s_%(process_num)02d 3command=php /home/forge/app.com/artisan queue:work --sleep=3 --tries=3 --max-time=3600 4autostart=true 5autorestart=true 6stopasgroup=true 7killasgroup=true 8user=forge 9numprocs=810redirect_stderr=true11stdout_logfile=/home/forge/app.com/worker.log12stopwaitsecs=3600
在此示例中,numprocs 指令将指示 Supervisor 运行八个 queue:work 进程并监视所有这些进程,如果它们失败则自动重启它们。你应该更改配置中的 command 指令以反映你所需的队列连接和工作程序选项。
你应该确保 stopwaitsecs 的值大于你运行时间最长的作业所消耗的秒数。否则,Supervisor 可能会在作业完成处理之前将其杀死。
启动 Supervisor
创建配置文件后,你可以使用以下命令更新 Supervisor 配置并启动进程:
1sudo supervisorctl reread2 3sudo supervisorctl update4 5sudo supervisorctl start "laravel-worker:*"
有关 Supervisor 的更多信息,请参阅 Supervisor 文档。
处理失败的任务
有时你的队列作业会失败。别担心,事情并不总是按计划进行!Laravel 包含了一种便捷的方法来指定作业应尝试的最大次数。当异步作业超过此尝试次数后,它将被插入到 failed_jobs 数据库表中。失败的同步分发作业不会存储在此表中,其异常由应用程序直接处理。
创建 failed_jobs 表的迁移文件通常已经存在于新的 Laravel 应用程序中。但是,如果你的应用程序不包含此表的迁移文件,你可以使用 make:queue-failed-table 命令来创建它:
1php artisan make:queue-failed-table2 3php artisan migrate
运行队列工作程序进程时,你可以使用 queue:work 命令上的 --tries 开关来指定作业应尝试的最大次数。如果你没有为 --tries 选项指定值,作业将只会尝试一次,或者按照作业类中 Tries 属性指定的次数进行尝试。
1php artisan queue:work redis --tries=3
使用 --backoff 选项,你可以指定 Laravel 在遇到异常后重试作业前应等待的秒数。默认情况下,作业会立即被释放回队列以便再次尝试:
1php artisan queue:work redis --tries=3 --backoff=3
如果你想针对每个作业配置 Laravel 在遇到异常后重试作业前应等待的秒数,可以在你的作业类上使用 Backoff 属性:
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\Backoff; 6 7#[Backoff(3)] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
如果你需要更复杂的逻辑来确定作业的退避(backoff)时间,可以在你的作业类上定义一个 backoff 方法:
1/**2 * Calculate the number of seconds to wait before retrying the job.3 */4public function backoff(): int5{6 return 3;7}
你可以通过定义退避值数组来轻松配置“指数”退避。在此示例中,第一次重试的延迟为 1 秒,第二次为 5 秒,第三次为 10 秒,如果还有剩余尝试次数,则随后的每次重试均为 10 秒:
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\Backoff; 6 7#[Backoff([1, 5, 10])] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
清理失败任务后的处理
当特定作业失败时,你可能希望向用户发送提醒或撤销该作业部分完成的操作。为此,你可以在你的作业类中定义一个 failed 方法。导致作业失败的 Throwable 实例将被传递给 failed 方法:
1<?php 2 3namespace App\Jobs; 4 5use App\Models\Podcast; 6use App\Services\AudioProcessor; 7use Illuminate\Contracts\Queue\ShouldQueue; 8use Illuminate\Foundation\Queue\Queueable; 9use Throwable;10 11class ProcessPodcast implements ShouldQueue12{13 use Queueable;14 15 /**16 * Create a new job instance.17 */18 public function __construct(19 public Podcast $podcast,20 ) {}21 22 /**23 * Execute the job.24 */25 public function handle(AudioProcessor $processor): void26 {27 // Process uploaded podcast...28 }29 30 /**31 * Handle a job failure.32 */33 public function failed(?Throwable $exception): void34 {35 // Send user notification of failure, etc...36 }37}
在调用 failed 方法之前,会实例化该作业的一个新实例;因此,在 handle 方法中可能发生的任何类属性修改都将丢失。
失败的作业不一定是指遇到了未捕获异常的作业。当作业用尽了所有允许的尝试次数时,它也可能被视为失败。这些尝试次数可以通过多种方式消耗:
- 任务超时。
- 任务在执行期间遇到未处理的异常。
- 作业被手动或通过中间件释放回队列。
如果最后一次尝试因作业执行期间抛出的异常而失败,该异常将被传递给作业的 failed 方法。但是,如果作业是因为达到了允许的最大尝试次数而失败,则 $exception 将是 Illuminate\Queue\MaxAttemptsExceededException 的实例。同样,如果作业因超过配置的超时时间而失败,则 $exception 将是 Illuminate\Queue\TimeoutExceededException 的实例。
重试失败任务
要查看已插入到 failed_jobs 数据库表中的所有失败作业,可以使用 queue:failed Artisan 命令:
1php artisan queue:failed
queue:failed 命令将列出作业 ID、连接、队列、失败时间以及有关该作业的其他信息。作业 ID 可用于重试失败的作业。例如,要重试 ID 为 ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece 的失败作业,请执行以下命令:
1php artisan queue:retry ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece
如有必要,你可以将多个 ID 传递给该命令:
1php artisan queue:retry ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece 91401d2c-0784-4f43-824c-34f94a33c24d
你还可以重试特定队列的所有失败作业:
1php artisan queue:retry --queue=name
要重试所有失败的作业,请执行 queue:retry 命令并将 all 作为 ID 传递:
1php artisan queue:retry all
如果你想删除一个失败的作业,可以使用 queue:forget 命令:
1php artisan queue:forget 91401d2c-0784-4f43-824c-34f94a33c24d
使用 Horizon 时,你应该使用 horizon:forget 命令而不是 queue:forget 命令来删除失败的作业。
要从 failed_jobs 表中删除所有失败的作业,可以使用 queue:flush 命令:
1php artisan queue:flush
queue:flush 命令会删除队列中的所有失败作业记录,无论失败作业是什么时候发生的。你可以使用 --hours 选项仅删除在一定小时数前或更早时间失败的作业:
1php artisan queue:flush --hours=48
忽略缺失模型
当将 Eloquent 模型注入到作业中时,模型会在放入队列之前自动序列化,并在作业被处理时从数据库中重新获取。但是,如果模型在作业等待工作程序处理时被删除,你的作业可能会以 ModelNotFoundException 失败。
为了方便起见,你可以选择使用作业类上的 DeleteWhenMissingModels 属性自动删除缺少模型的作业。当存在此属性时,Laravel 将静默丢弃该作业而不引发异常:
1<?php 2 3namespace App\Jobs; 4 5use Illuminate\Queue\Attributes\DeleteWhenMissingModels; 6 7#[DeleteWhenMissingModels] 8class ProcessPodcast implements ShouldQueue 9{10 // ...11}
清理失败任务
你可以通过调用 queue:prune-failed Artisan 命令来清理应用程序 failed_jobs 表中的记录:
1php artisan queue:prune-failed
默认情况下,所有超过 24 小时的失败作业记录都将被清理。如果你为命令提供了 --hours 选项,则仅保留过去 N 小时内插入的失败作业记录。例如,以下命令将删除所有在 48 小时前插入的失败作业记录:
1php artisan queue:prune-failed --hours=48
在 DynamoDB 中存储失败任务
Laravel 还支持将失败的作业记录存储在 DynamoDB 中,而不是关系数据库表中。但是,你必须手动创建一个 DynamoDB 表来存储所有失败的作业记录。通常,此表应命名为 failed_jobs,但你应该根据应用程序 queue 配置文件中 queue.failed.table 配置值来命名表。
failed_jobs 表应具有一个名为 application 的字符串主分区键和一个名为 uuid 的字符串主排序键。键中的 application 部分将包含你的应用程序名称(由应用程序 app 配置文件中的 name 配置值定义)。由于应用程序名称是 DynamoDB 表键的一部分,因此你可以使用同一个表来存储多个 Laravel 应用程序的失败作业。
此外,请确保安装了 AWS SDK,以便你的 Laravel 应用程序可以与 Amazon DynamoDB 通信:
1composer require aws/aws-sdk-php
接下来,将 queue.failed.driver 配置选项的值设置为 dynamodb。此外,你应该在失败作业配置数组中定义 key、secret 和 region 配置选项。这些选项将用于向 AWS 进行身份验证。使用 dynamodb 驱动程序时,queue.failed.database 配置选项是不必要的。
1'failed' => [2 'driver' => env('QUEUE_FAILED_DRIVER', 'dynamodb'),3 'key' => env('AWS_ACCESS_KEY_ID'),4 'secret' => env('AWS_SECRET_ACCESS_KEY'),5 'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),6 'table' => 'failed_jobs',7],
禁用失败任务存储
你可以通过将 queue.failed.driver 配置选项的值设置为 null,来指示 Laravel 丢弃失败的作业而不存储它们。通常,这可以通过 QUEUE_FAILED_DRIVER 环境变量来实现。
1QUEUE_FAILED_DRIVER=null
失败任务事件
如果你想注册一个在作业失败时调用的事件监听器,可以使用 Queue 门面的 failing 方法。例如,我们可以从 Laravel 自带的 AppServiceProvider 的 boot 方法中将闭包附加到此事件:
1<?php 2 3namespace App\Providers; 4 5use Illuminate\Support\Facades\Queue; 6use Illuminate\Support\ServiceProvider; 7use Illuminate\Queue\Events\JobFailed; 8 9class AppServiceProvider extends ServiceProvider10{11 /**12 * Register any application services.13 */14 public function register(): void15 {16 // ...17 }18 19 /**20 * Bootstrap any application services.21 */22 public function boot(): void23 {24 Queue::failing(function (JobFailed $event) {25 // $event->connectionName26 // $event->job27 // $event->exception28 });29 }30}
从队列中清除任务
使用 Horizon 时,你应该使用 horizon:clear 命令而不是 queue:clear 命令来清除队列中的作业。
如果你想从默认连接的默认队列中删除所有作业,可以使用 queue:clear Artisan 命令:
1php artisan queue:clear
你还可以提供 connection 参数和 queue 选项,以从特定的连接和队列中删除作业:
1php artisan queue:clear redis --queue=emails
从队列中清除作业仅适用于 SQS、Redis 和数据库队列驱动程序。此外,SQS 消息删除过程最多需要 60 秒,因此在你清除队列后的 60 秒内发送到 SQS 队列的作业也可能会被删除。
监控队列
如果你的队列突然接收到大量作业,它可能会不堪重负,导致作业完成的等待时间变长。如果需要,当队列作业计数超过指定阈值时,Laravel 可以向你发出警报。
首先,你应该安排 queue:monitor 命令每分钟运行一次。该命令接受你希望监视的队列名称以及你想要的作业计数阈值:
1php artisan queue:monitor redis:default,redis:deployments --max=100
仅安排此命令不足以触发通知来提醒你队列过载。当命令遇到作业计数超过阈值的队列时,将分发一个 Illuminate\Queue\Events\QueueBusy 事件。你可以在应用程序的 AppServiceProvider 中监听此事件,以便向你或你的开发团队发送通知:
1use App\Notifications\QueueHasLongWaitTime; 2use Illuminate\Queue\Events\QueueBusy; 3use Illuminate\Support\Facades\Event; 4use Illuminate\Support\Facades\Notification; 5 6/** 7 * Bootstrap any application services. 8 */ 9public function boot(): void10{11 Event::listen(function (QueueBusy $event) {13 ->notify(new QueueHasLongWaitTime(14 $event->connectionName,15 $event->queue,16 $event->size17 ));18 });19}
测试
在测试分发作业的代码时,你可能希望指示 Laravel 不要实际执行作业,因为作业的代码可以直接测试,并与分发它的代码分开测试。当然,要测试作业本身,你可以在测试中实例化一个作业实例并直接调用 handle 方法。
你可以使用 Queue 门面的 fake 方法来防止队列作业被实际推送到队列。调用 Queue 门面的 fake 方法后,你可以断言应用程序尝试将作业推送到队列:
1<?php 2 3use App\Jobs\AnotherJob; 4use App\Jobs\ShipOrder; 5use Illuminate\Support\Facades\Queue; 6 7test('orders can be shipped', function () { 8 Queue::fake(); 9 10 // Perform order shipping...11 12 // Assert that no jobs were pushed...13 Queue::assertNothingPushed();14 15 // Assert a job was pushed to a given queue...16 Queue::assertPushedOn('queue-name', ShipOrder::class);17 18 // Assert a job was pushed19 Queue::assertPushed(ShipOrder::class);20 21 // Assert a job was pushed twice...22 Queue::assertPushedTimes(ShipOrder::class, 2);23 24 // Assert a job was not pushed...25 Queue::assertNotPushed(AnotherJob::class);26 27 // Assert that a closure was pushed to the queue...28 Queue::assertClosurePushed();29 30 // Assert that a closure was not pushed...31 Queue::assertClosureNotPushed();32 33 // Assert the total number of jobs that were pushed...34 Queue::assertCount(3);35});
1<?php 2 3namespace Tests\Feature; 4 5use App\Jobs\AnotherJob; 6use App\Jobs\ShipOrder; 7use Illuminate\Support\Facades\Queue; 8use Tests\TestCase; 9 10class ExampleTest extends TestCase11{12 public function test_orders_can_be_shipped(): void13 {14 Queue::fake();15 16 // Perform order shipping...17 18 // Assert that no jobs were pushed...19 Queue::assertNothingPushed();20 21 // Assert a job was pushed to a given queue...22 Queue::assertPushedOn('queue-name', ShipOrder::class);23 24 // Assert a job was pushed25 Queue::assertPushed(ShipOrder::class);26 27 // Assert a job was pushed twice...28 Queue::assertPushedTimes(ShipOrder::class, 2);29 30 // Assert a job was not pushed...31 Queue::assertNotPushed(AnotherJob::class);32 33 // Assert that a closure was pushed to the queue...34 Queue::assertClosurePushed();35 36 // Assert that a closure was not pushed...37 Queue::assertClosureNotPushed();38 39 // Assert the total number of jobs that were pushed...40 Queue::assertCount(3);41 }42}
你可以将闭包传递给 assertPushed、assertNotPushed、assertClosurePushed 或 assertClosureNotPushed 方法,以断言推送了一个通过特定“真值测试”的作业。如果推送了至少一个通过该测试的作业,则断言将成功:
1use Illuminate\Queue\CallQueuedClosure;2 3Queue::assertPushed(function (ShipOrder $job) use ($order) {4 return $job->order->id === $order->id;5});6 7Queue::assertClosurePushed(function (CallQueuedClosure $job) {8 return $job->name === 'validate-order';9});
伪造部分任务
如果你只需要伪造特定的作业,同时允许其他作业正常执行,可以将应该伪造的作业类名传递给 fake 方法:
1test('orders can be shipped', function () { 2 Queue::fake([ 3 ShipOrder::class, 4 ]); 5 6 // Perform order shipping... 7 8 // Assert a job was pushed twice... 9 Queue::assertPushedTimes(ShipOrder::class, 2);10});
1public function test_orders_can_be_shipped(): void 2{ 3 Queue::fake([ 4 ShipOrder::class, 5 ]); 6 7 // Perform order shipping... 8 9 // Assert a job was pushed twice...10 Queue::assertPushedTimes(ShipOrder::class, 2);11}
你可以使用 except 方法伪造除一组指定作业之外的所有作业:
1Queue::fake()->except([2 ShipOrder::class,3]);
测试任务链
要测试作业链(Job Chaining),你需要利用 Bus 门面的伪造功能。Bus 门面的 assertChained 方法可用于断言一个作业链已被分发。assertChained 方法接受一个链式作业数组作为其第一个参数:
1use App\Jobs\RecordShipment; 2use App\Jobs\ShipOrder; 3use App\Jobs\UpdateInventory; 4use Illuminate\Support\Facades\Bus; 5 6Bus::fake(); 7 8// ... 9 10Bus::assertChained([11 ShipOrder::class,12 RecordShipment::class,13 UpdateInventory::class14]);
如上例所示,链式作业数组可以是一组作业的类名。但是,你也可以提供一组实际的作业实例。这样做时,Laravel 将确保作业实例属于同一类,并具有与应用程序分发的链式作业相同的属性值。
1Bus::assertChained([2 new ShipOrder,3 new RecordShipment,4 new UpdateInventory,5]);
你可以使用 assertDispatchedWithoutChain 方法来断言一个作业在没有作业链的情况下被推送。
1Bus::assertDispatchedWithoutChain(ShipOrder::class);
测试链式修改
如果链式作业向现有链中预置或追加了作业,你可以使用作业的 assertHasChain 方法来断言该作业具有预期的剩余作业链:
1$job = new ProcessPodcast;2 3$job->handle();4 5$job->assertHasChain([6 new TranscribePodcast,7 new OptimizePodcast,8 new ReleasePodcast,9]);
assertDoesntHaveChain 方法可用于断言作业的剩余链为空。
1$job->assertDoesntHaveChain();
测试链式批处理
如果你的作业链包含一组作业批处理,你可以通过在链断言中插入 Bus::chainedBatch 定义来断言该链式批处理符合你的预期:
1use App\Jobs\ShipOrder; 2use App\Jobs\UpdateInventory; 3use Illuminate\Bus\PendingBatch; 4use Illuminate\Support\Facades\Bus; 5 6Bus::assertChained([ 7 new ShipOrder, 8 Bus::chainedBatch(function (PendingBatch $batch) { 9 return $batch->jobs->count() === 3;10 }),11 new UpdateInventory,12]);
测试任务批处理
Bus 门面的 assertBatched 方法可用于断言一个作业批处理已被分发。传递给 assertBatched 方法的闭包接收一个 Illuminate\Bus\PendingBatch 实例,该实例可用于检查批处理中的作业:
1use Illuminate\Bus\PendingBatch; 2use Illuminate\Support\Facades\Bus; 3 4Bus::fake(); 5 6// ... 7 8Bus::assertBatched(function (PendingBatch $batch) { 9 return $batch->name == 'Import CSV' &&10 $batch->jobs->count() === 10;11});
hasJobs 方法可用于挂起的批处理,以验证批处理是否包含预期的作业。该方法接受一组作业实例、类名或闭包:
1Bus::assertBatched(function (PendingBatch $batch) {2 return $batch->hasJobs([3 new ProcessCsvRow(row: 1),4 new ProcessCsvRow(row: 2),5 new ProcessCsvRow(row: 3),6 ]);7});
使用闭包时,闭包将接收作业实例。预期的作业类型将从闭包的类型提示中推断出来。
1Bus::assertBatched(function (PendingBatch $batch) {2 return $batch->hasJobs([3 fn (ProcessCsvRow $job) => $job->row === 1,4 fn (ProcessCsvRow $job) => $job->row === 2,5 fn (ProcessCsvRow $job) => $job->row === 3,6 ]);7});
你可以使用 assertBatchCount 方法来断言分发了指定数量的批处理。
1Bus::assertBatchCount(3);
你可以使用 assertNothingBatched 来断言没有分发任何批处理。
1Bus::assertNothingBatched();
测试作业/批处理交互
此外,你有时可能需要测试单个作业与其底层批处理的交互。例如,你可能需要测试作业是否取消了其批处理的进一步处理。为此,你需要通过 withFakeBatch 方法为作业分配一个伪批处理。withFakeBatch 方法返回一个包含作业实例和伪批处理的元组:
1[$job, $batch] = (new ShipOrder)->withFakeBatch();2 3$job->handle();4 5$this->assertTrue($batch->cancelled());6$this->assertEmpty($batch->added);
测试任务/队列交互
有时,你可能需要测试队列作业是否将自己释放回队列。或者,你可能需要测试作业是否删除了自己。你可以通过实例化作业并调用 withFakeQueueInteractions 方法来测试这些队列交互。
一旦作业的队列交互被伪造,你就可以在作业上调用 handle 方法。调用作业后,可以使用各种断言方法来验证作业的队列交互:
1use App\Exceptions\CorruptedAudioException; 2use App\Jobs\ProcessPodcast; 3 4$job = (new ProcessPodcast)->withFakeQueueInteractions(); 5 6$job->handle(); 7 8$job->assertReleased(delay: 30); 9$job->assertDeleted();10$job->assertNotDeleted();11$job->assertFailed();12$job->assertFailedWith(CorruptedAudioException::class);13$job->assertNotFailed();
任务事件
使用 Queue 门面上的 before 和 after 方法,你可以指定在处理队列作业之前或之后执行的回调。这些回调是执行额外日志记录或为仪表板增加统计数据的绝佳机会。通常,你应该从服务提供者的 boot 方法中调用这些方法。例如,我们可以使用 Laravel 自带的 AppServiceProvider:
1<?php 2 3namespace App\Providers; 4 5use Illuminate\Support\Facades\Queue; 6use Illuminate\Support\ServiceProvider; 7use Illuminate\Queue\Events\JobProcessed; 8use Illuminate\Queue\Events\JobProcessing; 9 10class AppServiceProvider extends ServiceProvider11{12 /**13 * Register any application services.14 */15 public function register(): void16 {17 // ...18 }19 20 /**21 * Bootstrap any application services.22 */23 public function boot(): void24 {25 Queue::before(function (JobProcessing $event) {26 // $event->connectionName27 // $event->job28 // $event->job->payload()29 });30 31 Queue::after(function (JobProcessed $event) {32 // $event->connectionName33 // $event->job34 // $event->job->payload()35 });36 }37}
使用 Queue 门面上的 looping 方法,你可以指定在工作程序尝试从队列中获取作业之前执行的回调。例如,你可能注册一个闭包来回滚之前失败的作业留下的任何未关闭事务:
1use Illuminate\Support\Facades\DB;2use Illuminate\Support\Facades\Queue;3 4Queue::looping(function () {5 while (DB::transactionLevel() > 0) {6 DB::rollBack();7 }8});