Laravel 队列

简介

构建 Web 应用时,有些任务,例如解析并保存上传的 CSV 文件,所需时间可能太长,不适合在普通 Web 请求期间执行。Laravel 让你能够轻松创建可在后台处理的队列任务。将耗时任务移入队列后,应用可以极快地响应 Web 请求,为客户提供更好的使用体验。

Laravel 队列针对各种队列后端提供统一的队列 API,例如 Amazon SQS、Redis,甚至关系型数据库。

Laravel 的队列配置选项存储在应用的 config/queue.php 配置文件中。这里包含框架内置各队列驱动的连接配置,包括数据库、Amazon SQS、Redis 和 Beanstalkd 驱动,以及用于开发或测试、会立即执行任务的同步驱动。此外,还提供了一个会丢弃队列任务的 null 队列驱动。

Laravel Horizon 为基于 Redis 的队列提供美观的仪表板和配置系统。详情请参阅完整的 Horizon 文档。

连接与队列的区别

开始使用 Laravel 队列之前,应先理解“连接”与“队列”的区别。config/queue.php 配置文件中有一个 connections 配置数组,用来定义与 Amazon SQS、Beanstalk 或 Redis 等后端队列服务的连接。不过,一个队列连接可以包含多个“队列”,可以将它们理解为不同的任务集合,每个集合里都存放着等待处理的任务。

注意,queue 配置文件中的每个连接配置示例都包含一个 queue 属性。向特定连接发送任务时,该属性指定任务默认派发到哪个队列。也就是说,如果派发任务时没有明确指定队列,任务就会被放入该连接配置的 queue 属性所指定的队列:


use App\Jobs\ProcessPodcast;

// This job is sent to the default connection's default queue...
ProcessPodcast::dispatch();

// This job is sent to the default connection's "emails" queue...
ProcessPodcast::dispatch()->onQueue('emails');

有些应用可能永远不需要将任务推入多个队列,而只需一个简单队列。不过,对于希望按优先级或类别处理任务的应用,将任务推入多个队列尤其有用,因为 Laravel 队列工作进程允许你按优先级指定它应该处理哪些队列。例如,将任务推入 high 队列后,可以运行一个优先处理这些任务的工作进程:


php artisan queue:work --queue=high,default

驱动说明与前置条件

数据库

使用 database 队列驱动,需要一个存放任务的数据库表。通常,Laravel 默认的 0001_01_01_000002_create_jobs_table.php 数据库迁移 已包含这张表;如果应用没有该迁移,可以使用 make:queue-table Artisan 命令创建迁移:


php artisan make:queue-table

php artisan migrate

Redis

使用 redis 队列驱动时,应在 config/database.php 配置文件中配置 Redis 数据库连接。

Redis 的 serializer 和 compression 选项不受 redis 队列驱动支持。

Redis 集群

如果 Redis 队列连接使用 Redis 集群,队列名称必须包含一个 键哈希标签。这是为了确保同一队列的所有 Redis 键都被放在同一个哈希槽中:


'redis' => [
    'driver' => 'redis',
    'connection' => env('REDIS_QUEUE_CONNECTION', 'default'),
    'queue' => env('REDIS_QUEUE', '{default}'),
    'retry_after' => env('REDIS_QUEUE_RETRY_AFTER', 90),
    'block_for' => null,
    'after_commit' => false,
],
阻塞

使用 Redis 队列时,可以通过 block_for 配置选项指定驱动等待任务可用的时间,然后再进入工作进程循环的下一轮,并重新轮询 Redis 数据库。

根据队列负载调整这个值,比持续轮询 Redis 数据库查找新任务更高效。例如,可将其设为 5,表示驱动在等待任务可用时阻塞五秒:


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

将 block_for 设为 0,会让队列工作进程无限期阻塞,直到有任务可用。这也会使 SIGTERM 等信号在下一个任务处理完成之前无法得到处理。

SQS 超限载荷存储

Amazon SQS 限制队列消息载荷的最大大小。如果需要派发的任务载荷可能超过这一限制,可以配置 Laravel,将过大的 SQS 载荷存入缓存存储,然后通过 SQS 发送指向载荷的引用。要启用此功能,请在 SQS 队列连接配置中加入一个 overflow 数组:


'sqs' => [
    'driver' => 'sqs',
    'key' => env('AWS_ACCESS_KEY_ID'),
    'secret' => env('AWS_SECRET_ACCESS_KEY'),
    'prefix' => env('SQS_PREFIX', 'https://sqs.us-east-1.amazonaws.com/your-account-id'),
    'queue' => env('SQS_QUEUE', 'default'),
    'suffix' => env('SQS_SUFFIX'),
    'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),
    'after_commit' => false,
    'overflow' => [
        'enabled' => env('SQS_OVERFLOW_ENABLED', false),
        'store' => env('SQS_OVERFLOW_STORE'),
        'always' => false,
        'delete_after_processing' => true,
        'flush_on_clear' => env('SQS_OVERFLOW_FLUSH_ON_CLEAR', false),
    ],
],

启用超限载荷存储后,Laravel 会将大小至少为 1 MB 的载荷存入配置的缓存存储。如果 always 选项为 true,每个 SQS 载荷都会存入缓存,无论大小如何。任务处理时需要从缓存取回载荷,因此,应选择一个能够保留载荷直到工作进程处理任务的存储。默认情况下,任务成功处理并从 SQS 删除后,存储的载荷也会被删除。

如果 flush_on_clear 选项为 true,当 queue:clear 命令清空 SQS 队列时,也会清空配置的超限载荷缓存存储。清空缓存存储可能移除其中所有项目,因此,启用该选项时,应为 SQS 超限载荷存储配置专用的缓存存储。

其他驱动的前置条件

下列队列驱动需要相应的依赖项。可以通过 Composer 包管理器安装这些依赖:

  • Amazon SQS:aws/aws-sdk-php ~3.0
  • Beanstalkd:pda/pheanstalk ^7.0|^8.0
  • Redis:predis/predis ~3.0 或 phpredis PHP 扩展
  • MongoDB:mongodb/laravel-mongodb

创建任务

生成任务类

默认情况下,应用中所有可加入队列的任务都存放在 app/Jobs 目录中。如果 app/Jobs 目录不存在,运行 make:job Artisan 命令时会创建它:


php artisan make:job ProcessPodcast

生成的类会实现 Illuminate\Contracts\Queue\ShouldQueue 接口,告诉 Laravel,该任务应该推入队列并异步运行。

可以通过发布桩文件来自定义任务桩文件。

类结构

任务类非常简单,通常只包含一个 handle 方法,队列处理任务时会调用它。先来看一个任务类示例。假设我们管理一个播客发布服务,需要在发布之前处理上传的播客文件:


<?php

namespace App\Jobs;

use App\Models\Podcast;
use App\Services\AudioProcessor;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public Podcast $podcast,
    ) {}

    /**
     * Execute the job.
     */
    public function handle(AudioProcessor $processor): void
    {
        // Process uploaded podcast...
    }
}

注意,在本例中,我们可以将 Eloquent 模型 直接传入队列任务的构造函数。由于任务使用了 Queueable trait,在任务处理过程中,Eloquent 模型及已加载的关联都能妥善地进行序列化和反序列化。

如果队列任务的构造函数接收一个 Eloquent 模型,那么加入队列时,只会序列化该模型的标识符。实际处理任务时,队列系统会自动从数据库重新获取完整的模型实例及其已加载的关联。这种模型序列化方式,使发送给队列驱动的任务载荷小得多。

handle 方法依赖注入

队列处理任务时,会调用 handle 方法。注意,我们可以在任务的 handle 方法上为依赖项声明类型。Laravel 的服务容器会自动注入这些依赖项。

如果希望完全控制容器向 handle 方法注入依赖项的方式,可以使用容器的 bindMethod 方法。bindMethod 方法接收一个回调,回调会收到任务和容器。在回调中,你可以按需要调用 handle 方法。通常应在 boot 方法中调用这个方法,该方法位于 App\Providers\AppServiceProvider 服务提供者中:


use App\Jobs\ProcessPodcast;
use App\Services\AudioProcessor;
use Illuminate\Contracts\Foundation\Application;

$this->app->bindMethod([ProcessPodcast::class, 'handle'], function (ProcessPodcast $job, Application $app) {
    return $job->handle($app->make(AudioProcessor::class));
});

二进制数据,例如原始图像内容,在传入队列任务之前,应先经过 base64_encode 函数处理。否则,将任务放入队列时,可能无法正确将其序列化为 JSON。

队列任务中的关联

任务加入队列时,Eloquent 模型上所有已加载的关联也会被序列化,因此,序列化后的任务字符串有时会变得相当大。此外,任务反序列化并从数据库重新获取模型关联时,会获取这些关联的全部内容。在任务入队、模型序列化之前施加的任何关联约束,都不会在任务反序列化时重新应用。因此,如果只希望使用某个关联的部分数据,应在队列任务内部重新为该关联设置约束。

也可以在为属性赋值时,对模型调用 withoutRelations 方法,避免序列化关联。该方法会返回一个不包含已加载关联的模型实例:


/**
 * Create a new job instance.
 */
public function __construct(
    Podcast $podcast,
) {
    $this->podcast = $podcast->withoutRelations();
}

如果只需要移除特定关联,同时保留其他关联,可以使用 withoutRelation 方法:


$this->podcast = $podcast->withoutRelation('comments');

如果使用 PHP 构造函数属性提升,并希望表明某个 Eloquent 模型不应序列化其关联,可以使用 WithoutRelations 特性:


use Illuminate\Queue\Attributes\WithoutRelations;

/**
 * Create a new job instance.
 */
public function __construct(
    #[WithoutRelations]
    public Podcast $podcast,
) {}

为了方便,如果希望所有模型都不序列化关联,可以将 WithoutRelations 特性应用于整个类,而不必分别为每个模型添加该特性:


<?php

namespace App\Jobs;

use App\Models\DistributionPlatform;
use App\Models\Podcast;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Illuminate\Queue\Attributes\WithoutRelations;

#[WithoutRelations]
class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public Podcast $podcast,
        public DistributionPlatform $platform,
    ) {}
}

如果任务接收的是 Eloquent 模型集合或数组,而不是单个模型,任务反序列化并执行时,不会恢复集合内各模型的关联。这是为了防止处理大量模型的任务消耗过多资源。

唯一任务

唯一任务需要支持锁的缓存驱动。目前,memcached、redis、dynamodb、database、file 和 array 缓存驱动支持原子锁。

唯一任务约束不适用于批次中的任务。

有时,你可能希望确保任意时刻,队列里都只有某个任务的一个实例。可以在任务类上实现 ShouldBeUnique 接口来做到这一点。该接口不要求在类中定义任何额外方法:


<?php

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Contracts\Queue\ShouldBeUnique;

class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    // ...
}

上例中的 UpdateSearchIndex 任务是唯一任务。因此,如果队列中已有该任务的另一个实例,且尚未处理完成,就不会再次派发它。

某些情况下,你可能希望定义一个让任务具有唯一性的特定“键”,或者指定超时时间,超过该时间后任务不再保持唯一。为此,可以使用 UniqueFor 特性,并在任务类中定义 uniqueId 方法:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Contracts\Queue\ShouldBeUnique;
use Illuminate\Queue\Attributes\UniqueFor;

#[UniqueFor(3600)]
class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    /**
     * The product instance.
     *
     * @var \App\Models\Product
     */
    public $product;

    /**
     * Get the unique ID for the job.
     */
    public function uniqueId(): string
    {
        return $this->product->id;
    }
}

上例中的 UpdateSearchIndex 任务按产品 ID 保持唯一。因此,使用相同产品 ID 新派发的任务,会被忽略,直到现有任务处理完成。此外,如果现有任务在一小时内未被处理,唯一锁会释放,便可以向队列派发另一个具有相同唯一键的任务。

如果应用从多个 Web 服务器或容器派发任务,应确保所有服务器都连接同一台中央缓存服务器,让 Laravel 能够准确判断任务是否唯一。

保持任务唯一,直到开始处理

默认情况下,唯一任务会在处理完成,或所有重试尝试都失败后“解锁”。但有些情况下,你可能希望任务在即将开始处理时立即解锁。为此,任务应实现 ShouldBeUniqueUntilProcessing 契约,而不是 ShouldBeUnique 契约:


<?php

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Contracts\Queue\ShouldBeUniqueUntilProcessing;

class UpdateSearchIndex implements ShouldQueue, ShouldBeUniqueUntilProcessing
{
    // ...
}

唯一任务锁

在内部,派发 ShouldBeUnique 任务时,Laravel 会尝试获取一个锁,使用 uniqueId 作为锁键。如果该锁已被持有,就不会派发任务。任务处理完成,或所有重试尝试都失败后,会释放该锁。默认情况下,Laravel 使用默认缓存驱动获取此锁。如果希望使用其他驱动获取锁,可以定义一个 uniqueVia 方法,返回应该使用的缓存驱动:


use Illuminate\Contracts\Cache\Repository;
use Illuminate\Support\Facades\Cache;

class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    // ...

    /**
     * Get the cache driver for the unique job lock.
     */
    public function uniqueVia(): Repository
    {
        return Cache::driver('redis');
    }
}

如果只需要限制任务的并发处理,应改用 WithoutOverlapping 任务中间件。

防抖任务

有时,你可能希望确保在短时间内多次派发同一个任务时,只有最后一次派发的任务实际执行。可以为任务添加 DebounceFor 特性来实现这一点:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Illuminate\Queue\Attributes\DebounceFor;

#[DebounceFor(30)]
class UpdateSearchIndex implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(public int $productId)
    {
    }

    /**
     * Get the debounce ID for the job.
     */
    public function debounceId(): string
    {
        return (string) $this->productId;
    }
}

上例中,针对同一产品反复派发 UpdateSearchIndex,若多次派发发生在 30 秒内,该任务会进行防抖处理,只有最后一次派发的任务会运行。

如果希望为频繁重新派发的任务设置最长延后时间,可以提供 maxWait 参数,将其传给 DebounceFor 特性:


#[DebounceFor(30, maxWait: 120)]
class UpdateSearchIndex implements ShouldQueue
{
    use Queueable;

    // ...
}

可以在任务中定义 debounceVia 方法,自定义用于跟踪防抖状态的缓存存储:


use Illuminate\Contracts\Cache\Repository;
use Illuminate\Support\Facades\Cache;

public function debounceVia(): Repository
{
    return Cache::driver('redis');
}

如果一个防抖任务被更新的派发所取代,Laravel 会派发 Illuminate\Queue\Events\JobDebounced 事件,并从队列中移除被取代的任务。

防抖任务与唯一任务互斥。使用 DebounceFor 特性的任务,不应实现 ShouldBeUnique。

如果应用从多个 Web 服务器或容器派发防抖任务,应确保所有服务器都连接同一台中央缓存服务器。

加密任务

Laravel 可以通过加密确保任务数据的隐私性和完整性。只需为任务类添加 ShouldBeEncrypted 接口即可开始使用。添加接口后,Laravel 会在将任务推入队列之前自动加密它:


<?php

use Illuminate\Contracts\Queue\ShouldBeEncrypted;
use Illuminate\Contracts\Queue\ShouldQueue;

class UpdateSearchIndex implements ShouldQueue, ShouldBeEncrypted
{
    // ...
}

任务中间件

任务中间件允许你在队列任务的执行前后包装自定义逻辑,减少任务自身的重复样板代码。例如,下面的 handle 方法使用 Laravel 的 Redis 速率限制功能,确保每五秒只处理一个任务:


use Illuminate\Support\Facades\Redis;

/**
 * Execute the job.
 */
public function handle(): void
{
    Redis::throttle('key')->block(0)->allow(1)->every(5)->then(function () {
        info('Lock obtained...');

        // Handle job...
    }, function () {
        // Could not obtain lock...

        return $this->release(5);
    });
}

虽然这段代码有效,但 handle 方法掺杂了 Redis 速率限制逻辑,显得繁杂。此外,任何需要进行速率限制的其他任务,都必须重复这些逻辑。我们可以将速率限制逻辑定义为任务中间件,而不放在 handle 方法中:


<?php

namespace App\Jobs\Middleware;

use Closure;
use Illuminate\Support\Facades\Redis;

class RateLimited
{
    /**
     * Process the queued job.
     *
     * @param  \Closure(object): void  $next
     */
    public function handle(object $job, Closure $next): void
    {
        Redis::throttle('key')
            ->block(0)->allow(1)->every(5)
            ->then(function () use ($job, $next) {
                // Lock obtained...

                $next($job);
            }, function () use ($job) {
                // Could not obtain lock...

                $job->release(5);
            });
    }
}

如你所见,与路由中间件类似,任务中间件会接收正在处理的任务,以及一个应在继续处理任务时调用的回调。

可以使用 make:job-middleware Artisan 命令生成新的任务中间件类。创建中间件后,可以从任务的 middleware 方法返回这些中间件,将它们附加到任务上。通过 make:job Artisan 命令生成的任务骨架没有这个方法,因此,你需要手动将它添加到任务类:


use App\Jobs\Middleware\RateLimited;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [new RateLimited];
}

任务中间件还可以分配给可加入队列的事件监听器、可邮寄对象和通知。

速率限制

前面演示了如何编写自己的任务速率限制中间件,但 Laravel 实际上已经内置了可用于限制任务速率的中间件。与路由速率限制器类似,任务速率限制器通过 RateLimiter facade 的 for 方法定义。

例如,你可能希望允许用户每小时备份一次数据,但不对高级客户施加这一限制。为此,可以定义一个 RateLimiter,放在 boot 方法中,该方法位于 AppServiceProvider:


use Illuminate\Cache\RateLimiting\Limit;
use Illuminate\Support\Facades\RateLimiter;

/**
 * Bootstrap any application services.
 */
public function boot(): void
{
    RateLimiter::for('backups', function (object $job) {
        return $job->user->vipCustomer()
            ? Limit::none()
            : Limit::perHour(1)->by($job->user->id);
    });
}

上例定义的是每小时的速率限制,但也可以通过 perMinute 方法轻松定义基于分钟的速率限制。此外,可以向速率限制器的 by 方法传入任意值,不过,这个值通常用于按客户区分速率限制:


return Limit::perMinute(50)->by($job->user->id);

定义速率限制后,可以使用 Illuminate\Queue\Middleware\RateLimited 中间件,将速率限制器附加到任务上。每当任务超过速率限制,中间件会根据速率限制的时长,设置适当延迟,将任务释放回队列:


use Illuminate\Queue\Middleware\RateLimited;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [new RateLimited('backups')];
}

将受到速率限制的任务释放回队列,仍然会增加任务的 attempts 总数。你可能需要相应调整任务类上的 Tries 和 MaxExceptions 特性。也可以使用 retryUntil 方法,指定一段时间,超过该时间后就不应再尝试执行任务。

使用 releaseAfter 方法,还可以指定将任务释放回队列后,必须经过多少秒才能再次尝试执行:


/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new RateLimited('backups'))->releaseAfter(60)];
}

如果不希望任务受到速率限制后重试,可以使用 dontRelease 方法:


/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new RateLimited('backups'))->dontRelease()];
}

使用 Redis 进行速率限制

如果使用 Redis,可以使用 Illuminate\Queue\Middleware\RateLimitedWithRedis 中间件。它针对 Redis 做了专门优化,比基本的速率限制中间件更高效:


use Illuminate\Queue\Middleware\RateLimitedWithRedis;

public function middleware(): array
{
    return [new RateLimitedWithRedis('backups')];
}

可以使用 connection 方法指定中间件应使用哪个 Redis 连接:


return [(new RateLimitedWithRedis('backups'))->connection('limiter')];

防止任务重叠执行

Laravel 内置 Illuminate\Queue\Middleware\WithoutOverlapping 中间件,允许你基于任意键防止任务重叠执行。如果队列任务正在修改某个资源,而该资源一次只能由一个任务修改,这个中间件就很有用。

例如,假设有一个更新用户信用评分的队列任务,你希望避免同一用户 ID 的信用评分更新任务重叠执行。可以使用 WithoutOverlapping 中间件,从任务的 middleware 方法返回它即可:


use Illuminate\Queue\Middleware\WithoutOverlapping;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [new WithoutOverlapping($this->user->id)];
}

将重叠任务释放回队列,仍然会增加任务的总尝试次数。你可能需要相应调整任务类上的 Tries 和 MaxExceptions 特性。例如,如果保留 Tries 的默认值 1,就会导致所有重叠任务都无法稍后重试。

任何同类型的重叠任务都会被释放回队列。也可以指定将任务释放回队列后,必须经过多少秒才能再次尝试执行:


/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new WithoutOverlapping($this->order->id))->releaseAfter(60)];
}

如果希望立即删除所有重叠任务,让它们不再重试,可以使用 dontRelease 方法:


/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new WithoutOverlapping($this->order->id))->dontRelease()];
}

WithoutOverlapping 中间件基于 Laravel 的原子锁功能。有时,任务可能意外失败或超时,导致锁未释放。因此,可以使用 expireAfter 方法明确指定锁的过期时间。例如,下例会让 Laravel 在任务开始处理三分钟后释放 WithoutOverlapping 锁:


/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new WithoutOverlapping($this->order->id))->expireAfter(180)];
}

WithoutOverlapping 中间件需要支持锁的缓存驱动。目前,memcached、redis、dynamodb、database、file 和 array 缓存驱动支持原子锁。

在不同任务类之间共享锁键

默认情况下,WithoutOverlapping 中间件只会防止同一个类的任务重叠执行。因此,即使两个不同任务类使用相同的锁键,也不会阻止它们重叠执行。不过,可以使用 shared 方法,让 Laravel 在不同任务类之间共享这个键:


use Illuminate\Queue\Middleware\WithoutOverlapping;

class ProviderIsDown
{
    // ...

    public function middleware(): array
    {
        return [
            (new WithoutOverlapping("status:{$this->provider}"))->shared(),
        ];
    }
}

class ProviderIsUp
{
    // ...

    public function middleware(): array
    {
        return [
            (new WithoutOverlapping("status:{$this->provider}"))->shared(),
        ];
    }
}

异常节流

Laravel 内置 Illuminate\Queue\Middleware\ThrottlesExceptions 中间件,可以进行异常节流。任务抛出指定数量的异常后,后续所有执行尝试都会延迟,直到指定时间间隔结束。这种中间件尤其适合与不稳定的第三方服务交互的任务。

例如,假设一个与第三方 API 交互的队列任务开始抛出异常。可以使用 ThrottlesExceptions 中间件,从任务的 middleware 方法返回它,对异常进行节流。通常,这种中间件应与采用基于时间的尝试的任务搭配使用:


use DateTime;
use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [new ThrottlesExceptions(10, 5 * 60)];
}

/**
 * Determine the time at which the job should timeout.
 */
public function retryUntil(): DateTime
{
    return now()->plus(minutes: 30);
}

中间件构造函数接收的第一个参数,表示任务在触发节流之前可以抛出的异常数量;第二个参数表示触发节流后,要经过多少秒才能再次尝试执行任务。在上面的示例中,如果任务连续抛出 10 次异常,我们会等待 5 分钟后再次尝试,但仍受 30 分钟总时间限制约束。

任务抛出异常,但尚未达到异常阈值时,通常会立即重试。不过,在将中间件附加到任务时,可以调用 backoff 方法,指定这类任务应延迟的分钟数:


use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(10, 5 * 60))->backoff(5)];
}

backoff 方法还接收一个闭包,闭包会收到抛出的异常,让你能够动态决定延迟时间:


use App\Exceptions\RateLimitedException;
use Illuminate\Queue\Middleware\ThrottlesExceptions;
use Throwable;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(10, 5 * 60))->backoff(
        fn (Throwable $throwable) => $throwable instanceof RateLimitedException
            ? $throwable->retryAfterMinutes()
            : 5
    )];
}

在内部,该中间件使用 Laravel 的缓存系统实现速率限制,并将任务类名用作缓存“键”。将中间件附加到任务时,可以调用 by 方法覆盖这个键。如果多个任务都与同一个第三方服务交互,而你希望它们共用一个节流“桶”,遵守同一份共享限额,这个功能会很有用:


use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(10, 10 * 60))->by('key')];
}

默认情况下,中间件会对每个异常进行节流。将中间件附加到任务时,可以调用 when 方法修改这一行为。这样,只有提供给 when 方法的闭包返回 true 时,才会对该异常进行节流:


use Illuminate\Http\Client\HttpClientException;
use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(10, 10 * 60))->when(
        fn (Throwable $throwable) => $throwable instanceof HttpClientException
    )];
}

when 方法会将任务释放回队列,或者抛出异常;与之不同,deleteWhen 方法允许在发生指定异常时,直接删除整个任务:


use App\Exceptions\CustomerDeletedException;
use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(2, 10 * 60))->deleteWhen(CustomerDeletedException::class)];
}

如果希望将受到节流的异常报告给应用的异常处理器,可以在将中间件附加到任务时调用 report 方法。也可以向 report 方法提供一个闭包,这样,只有该闭包返回 true 时才会报告异常:


use Illuminate\Http\Client\HttpClientException;
use Illuminate\Queue\Middleware\ThrottlesExceptions;

/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [(new ThrottlesExceptions(10, 10 * 60))->report(
        fn (Throwable $throwable) => $throwable instanceof HttpClientException
    )];
}

使用 Redis 进行异常节流

如果使用 Redis,可以使用 Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis 中间件。它针对 Redis 做了专门优化,比基本的异常节流中间件更高效:


use Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis;

public function middleware(): array
{
    return [new ThrottlesExceptionsWithRedis(10, 10 * 60)];
}

可以使用 connection 方法指定中间件应使用哪个 Redis 连接:


return [(new ThrottlesExceptionsWithRedis(10, 10 * 60))->connection('limiter')];

将任务释放回队列

Release 中间件允许将任务释放回队列,而不执行任务。Release::when 方法会在给定条件求值为 true 时释放任务;Release::unless 方法则会在条件求值为 false 时释放任务:


use Illuminate\Queue\Middleware\Release;

/**
 * Get the middleware the job should pass through.
 */
public function middleware(): array
{
    return [
        Release::when($condition, releaseAfter: 60),
    ];
}

将任务释放回队列,仍然会增加任务的总尝试次数。你可能需要相应调整任务类上的 Tries 和 MaxExceptions 特性。

还可以传入一个 Closure,将其提供给 when 和 unless 方法,进行更复杂的条件判断:


use Illuminate\Queue\Middleware\Release;

/**
 * Get the middleware the job should pass through.
 */
public function middleware(): array
{
    return [
        Release::when(function (): bool {
            return ! $this->order->isPaid();
        }, releaseAfter: 60),
    ];
}

跳过任务

Skip 中间件允许指定某个任务应直接跳过或删除,而无需修改任务逻辑。Skip::when 方法会在给定条件求值为 true 时删除任务;Skip::unless 方法则会在条件求值为 false 时删除任务:


use Illuminate\Queue\Middleware\Skip;

/**
 * Get the middleware the job should pass through.
 */
public function middleware(): array
{
    return [
        Skip::when($condition),
    ];
}

还可以传入一个 Closure,将其提供给 when 和 unless 方法,进行更复杂的条件判断:


use Illuminate\Queue\Middleware\Skip;

/**
 * Get the middleware the job should pass through.
 */
public function middleware(): array
{
    return [
        Skip::when(function (): bool {
            return $this->shouldSkip();
        }),
    ];
}

派发任务

编写好任务类后,可以通过任务自身的 dispatch 方法派发它。传给 dispatch 方法的参数会被传入任务构造函数:


<?php

namespace App\Http\Controllers;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\RedirectResponse;
use Illuminate\Http\Request;

class PodcastController extends Controller
{
    /**
     * Store a new podcast.
     */
    public function store(Request $request): RedirectResponse
    {
        $podcast = Podcast::create(/* ... */);

        // ...

        ProcessPodcast::dispatch($podcast);

        return redirect('/podcasts');
    }
}

如果希望根据条件派发任务,可以使用 dispatchIf 和 dispatchUnless 方法:


ProcessPodcast::dispatchIf($accountActive, $podcast);

ProcessPodcast::dispatchUnless($accountSuspended, $podcast);

在新建的 Laravel 应用中,database 连接被定义为默认队列。可以修改 QUEUE_CONNECTION 环境变量来指定其他默认队列连接,该变量位于应用的 .env 文件中。

延迟派发

如果希望任务派发后不能立即由队列工作进程处理,可以在派发时使用 delay 方法。例如,下面指定任务只有在派发十分钟后才能被处理:


<?php

namespace App\Http\Controllers;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\RedirectResponse;
use Illuminate\Http\Request;

class PodcastController extends Controller
{
    /**
     * Store a new podcast.
     */
    public function store(Request $request): RedirectResponse
    {
        $podcast = Podcast::create(/* ... */);

        // ...

        ProcessPodcast::dispatch($podcast)
            ->delay(now()->plus(minutes: 10));

        return redirect('/podcasts');
    }
}

某些情况下,任务可能配置了默认延迟。如果需要绕过这个延迟,并派发任务供立即处理,可以使用 withoutDelay 方法:


ProcessPodcast::dispatch($podcast)->withoutDelay();

Amazon SQS 队列服务的最大延迟时间为 15 分钟。

同步派发

如果希望立即(同步)派发任务,可以使用 dispatchSync 方法。使用此方法时,任务不会加入队列,而是在当前进程中立即执行:


<?php

namespace App\Http\Controllers;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\RedirectResponse;
use Illuminate\Http\Request;

class PodcastController extends Controller
{
    /**
     * Store a new podcast.
     */
    public function store(Request $request): RedirectResponse
    {
        $podcast = Podcast::create(/* ... */);

        // Create podcast...

        ProcessPodcast::dispatchSync($podcast);

        return redirect('/podcasts');
    }
}

延后派发

使用延后的同步派发,可以让任务仍在当前进程内处理,但等到 HTTP 响应已发送给用户后才执行。这样便能同步处理“队列”任务,而不降低用户使用应用时的体验。要延后执行同步任务,请将任务派发到 deferred 连接:


RecordDelivery::dispatch($order)->onConnection('deferred');

deferred 连接也充当默认的故障转移队列。

类似地,background 连接也会在 HTTP 响应已发送给用户后处理任务;不过,任务在一个单独启动的 PHP 进程中处理,让 PHP-FPM/应用工作进程能够处理另一个到来的 HTTP 请求:


RecordDelivery::dispatch($order)->onConnection('background');

批量派发

如果需要一次派发多个相互独立的任务,且不需要批次跟踪或回调,可以使用 bulk 方法,该方法由 Bus facade 提供。Laravel 会按任务配置的队列连接及队列名称分组,然后将每组任务批量推入相应队列:


use App\Jobs\ProcessUser;
use Illuminate\Support\Facades\Bus;

Bus::bulk(
    $users->map(fn ($user) => new ProcessUser($user))
);

派发前准备任务

如果任务需要在推入队列之前准备或检查自身状态,可以实现 Illuminate\Contracts\Queue\PreparesForDispatch 接口。Laravel 会在派发任务之前调用任务的 prepareForDispatch 方法。如果该方法返回 false,就不会派发任务:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\PreparesForDispatch;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Illuminate\Support\Facades\Cache;

class SyncPodcasts implements PreparesForDispatch, ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public array $podcastIds,
    ) {}

    /**
     * Prepare the job before dispatching.
     */
    public function prepareForDispatch(): bool
    {
        return collect($this->podcastIds)
            ->reject(fn (int $id) => Cache::has("podcast-syncing:{$id}"))
            ->isNotEmpty();
    }
}

任务与数据库事务

虽然在数据库事务中分发任务完全可行,但你应格外注意,确保任务确实能够成功执行。在事务中分发任务时,工作进程可能在外层事务提交之前就开始处理该任务。如果发生这种情况,你在数据库事务中对模型或数据库记录所做的更新,可能尚未体现在数据库中。此外,在事务中创建的模型或数据库记录可能还不存在于数据库中。

好在 Laravel 提供了几种方法来解决这个问题。首先,你可以在队列连接的配置数组中设置 after_commit 连接选项:


'redis' => [
    'driver' => 'redis',
    // ...
    'after_commit' => true,
],

当 after_commit 选项为 true 时,你仍可以在数据库事务中分发任务;不过,Laravel 会等待尚未结束的外层数据库事务提交之后,才真正分发任务。当然,如果当前没有尚未结束的数据库事务,任务就会立即分发。

如果事务因执行期间发生异常而回滚,在该事务中分发的任务将被丢弃。

将 after_commit 配置选项设为 true,还会使所有进入队列的事件监听器、邮件对象、通知和广播事件,在所有尚未结束的数据库事务提交之后再分发。

在分发时指定与事务提交相关的行为

即使没有将队列连接的 after_commit 配置选项设为 true,你仍然可以指定某个任务应在所有尚未结束的数据库事务提交之后再分发。为此,可以在分发操作后链式调用 afterCommit 方法:


use App\Jobs\ProcessPodcast;

ProcessPodcast::dispatch($podcast)->afterCommit();

同样,如果 after_commit 配置选项已设为 true,你也可以指定某个任务立即分发,而不等待任何尚未结束的数据库事务提交:


ProcessPodcast::dispatch($podcast)->beforeCommit();

任务链

任务链允许你指定一组队列任务,在首个任务成功执行后按顺序运行。如果序列中的某个任务失败,其余任务就不会运行。要执行队列任务链,可以使用 chain 方法,该方法由 Bus 门面提供。Laravel 的命令总线是一个底层组件,队列任务的分发功能建立在它之上:


use App\Jobs\OptimizePodcast;
use App\Jobs\ProcessPodcast;
use App\Jobs\ReleasePodcast;
use Illuminate\Support\Facades\Bus;

Bus::chain([
    new ProcessPodcast,
    new OptimizePodcast,
    new ReleasePodcast,
])->dispatch();

除了将任务类的实例加入链中,你还可以将闭包加入链中:


Bus::chain([
    new ProcessPodcast,
    new OptimizePodcast,
    function () {
        Podcast::update(/* ... */);
    },
])->dispatch();

在任务内部使用 $this->delete() 方法删除任务,并不会阻止后续链式任务被处理。只有任务链中的某个任务失败时,整个链才会停止执行。

任务链的连接与队列

如果想指定链式任务所使用的连接和队列,可以使用 onConnection 和 onQueue 方法。这些方法会指定应使用的队列连接和队列名称,除非某个队列任务被显式指定了不同的连接或队列:


Bus::chain([
    new ProcessPodcast,
    new OptimizePodcast,
    new ReleasePodcast,
])->onConnection('redis')->onQueue('podcasts')->dispatch();

向任务链添加任务

有时,你可能需要在任务链中的某个任务内部,将另一个任务添加到现有任务链的开头或末尾。可以使用 prependToChain 和 appendToChain 方法实现:


/**
 * Execute the job.
 */
public function handle(): void
{
    // ...

    // Prepend to the current chain, run job immediately after current job...
    $this->prependToChain(new TranscribePodcast);

    // Append to the current chain, run job at end of chain...
    $this->appendToChain(new TranscribePodcast);
}

任务链失败

构建任务链时,可以使用 catch 方法指定一个闭包,在链中的任务失败时调用。指定的回调会接收到导致任务失败的 Throwable 实例:


use Illuminate\Support\Facades\Bus;
use Throwable;

Bus::chain([
    new ProcessPodcast,
    new OptimizePodcast,
    new ReleasePodcast,
])->catch(function (Throwable $e) {
    // A job within the chain has failed...
})->dispatch();

由于任务链回调会被序列化,并在稍后由 Laravel 队列执行,因此不应在这些回调中使用 $this 变量。

自定义队列与连接

分发到指定队列

将任务推送到不同队列,可以对队列任务进行“分类”,甚至可以按优先级决定为各个队列分配多少工作进程。请注意,这并不会将任务推送到队列配置文件所定义的不同队列“连接”,而只是将任务推送到同一个连接中的特定队列。要指定队列,请在分发任务时使用 onQueue 方法:


<?php

namespace App\Http\Controllers;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\RedirectResponse;
use Illuminate\Http\Request;

class PodcastController extends Controller
{
    /**
     * Store a new podcast.
     */
    public function store(Request $request): RedirectResponse
    {
        $podcast = Podcast::create(/* ... */);

        // Create podcast...

        ProcessPodcast::dispatch($podcast)->onQueue('processing');

        return redirect('/podcasts');
    }
}

或者,你也可以在任务的构造函数中调用 onQueue 方法,指定任务所使用的队列:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct()
    {
        $this->onQueue('processing');
    }
}

分发到指定连接

如果应用会使用多个队列连接,可以通过 onConnection 方法指定将任务推送到哪个连接:


<?php

namespace App\Http\Controllers;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\RedirectResponse;
use Illuminate\Http\Request;

class PodcastController extends Controller
{
    /**
     * Store a new podcast.
     */
    public function store(Request $request): RedirectResponse
    {
        $podcast = Podcast::create(/* ... */);

        // Create podcast...

        ProcessPodcast::dispatch($podcast)->onConnection('sqs');

        return redirect('/podcasts');
    }
}

你可以链式调用 onConnection 和 onQueue 方法,为任务同时指定连接和队列:


ProcessPodcast::dispatch($podcast)
    ->onConnection('sqs')
    ->onQueue('processing');

或者,你也可以在任务的构造函数中调用 onConnection 方法,指定任务所使用的连接:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct()
    {
        $this->onConnection('sqs');
    }
}

队列路由

你可以使用 Queue 门面的 route 方法,为特定任务类定义默认连接和队列。如果希望某些任务始终使用特定队列,而不必在任务中指定连接或队列,这个方法就很有用。

除了为特定任务类设置路由,你还可以将接口、trait 或父类传给 route 方法。这样,任何实现该接口、使用该 trait 或继承该父类的任务,都会自动使用配置的连接和队列。

通常,应调用 route 方法,并将该调用放在服务提供者的 boot 方法中:


use App\Concerns\RequiresVideo;
use App\Jobs\ProcessPodcast;
use App\Jobs\ProcessVideo;
use Illuminate\Contracts\Broadcasting\ShouldBroadcast;
use Illuminate\Support\Facades\Queue;

/**
 * Bootstrap any application services.
 */
public function boot(): void
{
    Queue::route(ProcessPodcast::class, connection: 'redis', queue: 'podcasts');
    Queue::route(RequiresVideo::class, queue: 'video');
    Queue::route(ShouldBroadcast::class, queue: 'events');
}

如果只指定连接而没有指定队列,任务将被发送到默认队列:


Queue::route(ProcessPodcast::class, connection: 'redis');

还可以向 route 方法传入数组,一次为多个任务类设置路由:


Queue::route([
    ProcessPodcast::class => ['redis', 'podcasts'], // Connection and queue
    ProcessVideo::class => 'videos', // Queue only (uses default connection)
]);

任务自身仍然可以逐个覆盖队列路由配置。

你可以使用 forward 方法,将任务从一个队列转发到另一个队列或连接,也可以同时改变队列和连接。当你需要调整队列基础设施,而不想修改各个任务或分发位置时,这个方法很有用:


Queue::forward('reports', 'reports.fifo', 'sqs');
Queue::forward('payments', connection: 'sqs');
Queue::forward('updates', 'notifications');

还可以传入数组,一次转发多个队列:


Queue::forward([
    'reports' => 'reports.fifo',
    'emails' => 'emails.fifo',
], connection: 'sqs');

任务上显式配置的连接,优先于转发规则指定的连接。

指定任务的最大尝试次数与超时时间

最大尝试次数

任务尝试次数是 Laravel 队列系统的核心概念,也是许多高级功能的基础。初看时可能有些难理解,但在修改默认配置之前,应先了解它的工作方式。

任务被分发后,会被推送到队列中。随后,工作进程取出任务并尝试执行。这就是一次任务尝试。

不过,一次尝试并不一定意味着任务的 handle 方法已经执行。尝试次数也可能因以下几种情况而被“消耗”:

  • 任务在执行过程中遇到未处理的异常。
  • 使用 $this->release() 手动将任务释放回队列。
  • WithoutOverlapping 或 RateLimited 等中间件未能取得锁,从而释放任务。
  • 任务超时。
  • 任务的 handle 方法正常运行并完成,没有抛出异常。

你可能不希望无限期地尝试执行某个任务。因此,Laravel 提供了多种方式,用来指定任务可以尝试多少次,或可以在多长时间内进行尝试。

默认情况下,Laravel 只会尝试执行任务一次。如果任务使用了 WithoutOverlapping 或 RateLimited 等中间件,或者你会手动释放任务,通常需要通过 tries 选项增加允许的尝试次数。

指定任务最大尝试次数的一种方法,是使用 Artisan 命令行中的 --tries 选项。这会应用于该工作进程处理的所有任务,除非正在处理的任务自身已经指定了允许的尝试次数:


php artisan queue:work --tries=3

如果任务超过最大尝试次数,就会被视为“失败”任务。处理失败任务的更多信息,请参阅失败任务文档。如果将 --tries=0 传给 queue:work 命令,任务就会无限重试。

你也可以采取更细粒度的方式,在任务类本身使用 Tries 特性定义任务的最大尝试次数。如果任务上指定了最大尝试次数,它将优先于命令行提供的 --tries 值:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Tries;

#[Tries(5)]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

如果需要动态控制某个任务的最大尝试次数,可以在任务上定义 tries 方法:


/**
 * Determine number of times the job may be attempted.
 */
public function tries(): int
{
    return 5;
}

基于时间的尝试限制

除了定义任务失败前允许尝试多少次,你还可以定义一个时间,超过该时间就不再尝试执行任务。这样,任务就可以在给定时间范围内尝试任意次数。要定义停止尝试的时间,请在任务类中添加 retryUntil 方法。该方法应返回一个 DateTime 实例:


use DateTime;

/**
 * Determine the time at which the job should timeout.
 */
public function retryUntil(): DateTime
{
    return now()->plus(minutes: 10);
}

如果同时定义了 retryUntil 和 tries,Laravel 会优先使用 retryUntil 方法。

你也可以定义 Tries 特性或 retryUntil 方法,并将其用于队列事件监听器和队列通知。

最大异常次数

有时,你希望允许任务尝试很多次,但如果重试是由一定数量的未处理异常触发的,就应让任务失败;这与直接通过 release 方法释放任务的情况不同。为此,可以在任务类上使用 Tries 和 MaxExceptions 特性:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Illuminate\Queue\Attributes\MaxExceptions;
use Illuminate\Queue\Attributes\Tries;
use Illuminate\Support\Facades\Redis;

#[Tries(25)]
#[MaxExceptions(3)]
class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Execute the job.
     */
    public function handle(): void
    {
        Redis::throttle('key')->allow(10)->every(60)->then(function () {
            // Lock obtained, process the podcast...
        }, function () {
            // Unable to obtain lock...
            return $this->release(10);
        });
    }
}

在这个示例中,如果应用无法取得 Redis 锁,任务会被释放,并等待十秒后再执行,最多持续重试 25 次。但是,如果任务抛出了三个未处理的异常,任务就会失败。

默认情况下,如果某次尝试因工作进程崩溃或被终止而结束,例如内存耗尽,这次尝试不会计入任务的最大异常次数。如果希望将这些尝试也算作异常,可以在任务类上添加 CountCrashesAsExceptions 特性:


use Illuminate\Queue\Attributes\CountCrashesAsExceptions;
use Illuminate\Queue\Attributes\MaxExceptions;
use Illuminate\Queue\Attributes\Tries;

#[Tries(25)]
#[MaxExceptions(3)]
#[CountCrashesAsExceptions]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

使用该特性时,工作进程会在处理任务期间向应用缓存中保存一个标记。如果下次尝试执行任务时,该标记仍然存在,上一次尝试就会计为一次异常。

根据异常停止重试

有时,异常意味着队列任务应立即失败,而不是被释放后再次尝试。你可以使用异常配置的 dontRetry 方法,在应用的 bootstrap/app.php 文件中指定哪些异常类型应停止任务重试:


use App\Exceptions\InvalidPodcastSourceException;
use Illuminate\Foundation\Configuration\Exceptions;

->withExceptions(function (Exceptions $exceptions): void {
    $exceptions->dontRetry([
        InvalidPodcastSourceException::class,
    ]);
})

如果需要更细致地控制何时停止重试,可以向 dontRetryWhen 方法传入闭包。当闭包返回 true 时,任务会被标记为失败,不再重试:


use App\Exceptions\PodcastProcessingException;
use Illuminate\Foundation\Configuration\Exceptions;

->withExceptions(function (Exceptions $exceptions): void {
    $exceptions->dontRetryWhen(function (PodcastProcessingException $e) {
        return $e->reason() === 'Subscription expired';
    });
})

超时

通常,你大致知道队列任务预期需要多长时间。因此,Laravel 允许你指定“超时”值。默认超时时间为 60 秒。如果任务的处理时间超过超时值指定的秒数,处理该任务的工作进程就会因错误而退出。通常,工作进程会由服务器上配置的进程管理器自动重启。

可以使用 Artisan 命令行中的 --timeout 选项,指定任务最多允许运行多少秒:


php artisan queue:work --timeout=30

如果任务因持续超时而超过最大尝试次数,它就会被标记为失败。

你也可以在任务类上使用 Timeout 特性,定义任务最多允许运行多少秒。如果任务上指定了超时时间,它将优先于命令行指定的任何超时时间:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Timeout;

#[Timeout(120)]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

有时,套接字或对外 HTTP 连接等会阻塞 I/O 的操作,可能不会遵守你指定的超时时间。因此,使用这些功能时,应始终尽量通过其自身的 API 另行设置超时。例如,使用 Guzzle 时,应始终指定连接超时和请求超时。

要指定任务超时时间,必须安装 PCNTL PHP 扩展。此外,任务的“超时”值应始终小于其“retry after” 值。否则,任务可能在实际执行完毕或超时之前,就被再次尝试执行。--timeout 选项在调用 queue:work 命令并同时使用 --once 选项时不会生效。

超时时让任务失败

如果希望任务在超时时被标记为失败,可以在任务类上使用 FailOnTimeout 特性:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\FailOnTimeout;

#[FailOnTimeout]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

默认情况下,任务超时会消耗一次尝试次数,并被释放回队列(前提是允许重试)。但是,如果将任务配置为超时时失败,它就不会再重试,无论 tries 设置为何值。

SQS FIFO 队列与公平队列

Laravel 支持 Amazon SQS FIFO(先进先出)队列和公平队列。FIFO 队列允许严格按任务发送的顺序处理任务,并通过消息去重确保每条消息只被处理一次。

FIFO 队列需要消息组 ID,以确定哪些任务可以并行处理。组 ID 相同的任务会按顺序处理,而组 ID 不同的消息可以并发处理。

Laravel 提供了可链式调用的 onGroup 方法,用于在分发任务时指定消息组 ID:


ProcessOrder::dispatch($order)
    ->onGroup("customer-{$order->customer_id}");

如果分发到 SQS FIFO 队列的任务没有指定消息组,Laravel 会使用队列名称作为消息组 ID。

SQS FIFO 队列支持消息去重,以确保每条消息只被处理一次。在任务类中实现 deduplicationId 方法,即可提供自定义去重 ID:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ProcessSubscriptionRenewal implements ShouldQueue
{
    use Queueable;

    // ...

    /**
     * Get the job's deduplication ID.
     */
    public function deduplicationId(): string
    {
        return "renewal-{$this->subscription->id}";
    }
}

公平队列

如果使用 SQS 标准队列,设置消息组会启用公平排队。换句话说,一旦分配了消息组,SQS 就会利用这些分组,在不同租户或工作负载之间保持公平投递。无须进行额外的 Laravel 配置。

除了在分发时调用 onGroup,你也可以直接在任务上定义 messageGroup 方法:


<?php

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ProcessOrder implements ShouldQueue
{
    use Queueable;

    // ...

    /**
     * Get the job's message group.
     */
    public function messageGroup(): string
    {
        return "customer-{$this->order->customer_id}";
    }
}

FIFO 监听器、邮件与通知

使用 FIFO 队列时,你还需要为监听器、邮件和通知定义消息组。或者,也可以将这些对象的队列实例分发到非 FIFO 队列。

要为队列事件监听器定义消息组,请在监听器上定义 messageGroup 方法。还可以选择定义 deduplicator 方法,该方法接收事件,并应返回一个用于生成去重 ID 的闭包:


<?php

namespace App\Listeners;

use App\Events\OrderShipped;
use Closure;

class SendShipmentNotification
{
    // ...

    /**
     * Get the job's message group.
     */
    public function messageGroup(): string
    {
        return 'shipments';
    }

    /**
     * Get the job's deduplicator.
     */
    public function deduplicator(OrderShipped $event): Closure
    {
        return fn () => "shipment-notification-{$event->order->id}";
    }
}

发送将进入 FIFO 队列的邮件消息时,应在发送通知时调用 onGroup 方法,并可选择调用 withDeduplicator 方法:


use App\Mail\InvoicePaid;
use Illuminate\Support\Facades\Mail;

$invoicePaid = (new InvoicePaid($invoice))
    ->onGroup('invoices')
    ->withDeduplicator(fn () => 'invoices-'.$invoice->id);

Mail::to($request->user())->send($invoicePaid);

发送将进入 FIFO 队列的通知时,应在发送通知时调用 onGroup 方法,并可选择调用 withDeduplicator 方法:


use App\Notifications\InvoicePaid;

$invoicePaid = (new InvoicePaid($invoice))
    ->onGroup('invoices')
    ->withDeduplicator(fn () => 'invoices-'.$invoice->id);

$user->notify($invoicePaid);

队列故障转移

failover 队列驱动在将任务推送到队列时提供自动故障转移功能。如果 failover 配置中的首选队列连接因任何原因失败,Laravel 就会自动尝试将任务推送到列表中下一个配置的连接。这对于保证生产环境的高可用性特别有用,尤其是在队列可靠性至关重要的场景中。

要配置故障转移队列连接,请指定 failover 驱动,并提供一个按顺序尝试的连接名称数组。默认情况下,Laravel 在应用的 config/queue.php 配置文件中提供了一个故障转移配置示例:


'failover' => [
    'driver' => 'failover',
    'connections' => [
        'redis',
        'database',
        'sync',
    ],
],

配置好使用 failover 驱动的连接后,需要在应用的 .env 文件中将该故障转移连接设为默认队列连接,才能使用故障转移功能:


QUEUE_CONNECTION=failover

接下来,为故障转移连接列表中的每个连接启动至少一个工作进程:


php artisan queue:work redis
php artisan queue:work database

对于使用 sync、background 或 deferred 队列驱动的连接,无须运行工作进程,因为这些驱动会在当前 PHP 进程中处理任务。

当队列连接操作失败并触发故障转移时,Laravel 会分发 Illuminate\Queue\Events\QueueFailedOver 事件,让你可以报告或记录队列连接失败的情况。

如果使用 Laravel Horizon,请记住 Horizon 只管理 Redis 队列。如果故障转移列表中包含 database,应在运行 Horizon 的同时,运行一个常规的 php artisan queue:work database 进程。

错误处理

如果处理任务时抛出异常,任务会自动被释放回队列,以便再次尝试执行。任务会不断被释放,直到达到应用允许的最大尝试次数。最大尝试次数由 --tries 选项定义,该选项用于 queue:work Artisan 命令。或者,也可以直接在任务类上定义最大尝试次数。关于运行队列工作进程的更多信息,见下文。

手动释放任务

有时,你可能希望手动将任务释放回队列,以便稍后再次尝试。可以通过调用 release 方法实现:


/**
 * Execute the job.
 */
public function handle(): void
{
    // ...

    $this->release();
}

默认情况下,release 方法会将任务释放回队列,以便立即处理。但是,你也可以向 release 方法传入整数或日期实例,让队列在指定秒数过去之后,才允许再次处理该任务:


$this->release(10);

$this->release(now()->plus(seconds: 10));

手动让任务失败

有时,你可能需要手动将任务标记为“失败”。为此,可以调用 fail 方法:


/**
 * Execute the job.
 */
public function handle(): void
{
    // ...

    $this->fail();
}

如果想因为捕获到的异常而将任务标记为失败,可以将该异常传给 fail 方法。为方便起见,也可以传入一个错误消息字符串,它会自动转换为异常:


$this->fail($exception);

$this->fail('Something went wrong.');

有关失败任务的更多信息,请参阅处理任务失败的文档。

遇到特定异常时让任务失败

FailOnException 任务中间件允许在抛出特定异常时直接停止重试。这样,面对外部 API 错误等暂时性异常时可以重试;而面对用户权限被撤销等持续性异常时,则让任务永久失败:


<?php

namespace App\Jobs;

use App\Models\User;
use Illuminate\Auth\Access\AuthorizationException;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Illuminate\Queue\Attributes\Tries;
use Illuminate\Queue\Middleware\FailOnException;
use Illuminate\Support\Facades\Http;

#[Tries(3)]
class SyncChatHistory implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public User $user,
    ) {}

    /**
     * Execute the job.
     */
    public function handle(): void
    {
        $this->user->authorize('sync-chat-history');

        $response = Http::throw()->get(
            "https://chat.laravel.test/?user={$this->user->uuid}"
        );

        // ...
    }

    /**
     * Get the middleware the job should pass through.
     */
    public function middleware(): array
    {
        return [
            new FailOnException([AuthorizationException::class])
        ];
    }
}

任务批处理

Laravel 的任务批处理功能让你可以轻松地并行执行一组任务,并在这批任务全部执行完毕后执行某些操作。

开始之前,应创建一个数据库迁移,建立用于保存任务批次元信息的表,例如批次的完成百分比。可以使用 make:queue-batches-table Artisan 命令生成该迁移:


php artisan make:queue-batches-table

php artisan migrate

定义可批处理的任务

要定义可批处理的任务,应照常创建可加入队列的任务,但需要向任务类添加 Illuminate\Bus\Batchable trait。这个 trait 提供了 batch 方法,可用于获取任务当前所在的批次:


<?php

namespace App\Jobs;

use Illuminate\Bus\Batchable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ImportCsv implements ShouldQueue
{
    use Batchable, Queueable;

    /**
     * Execute the job.
     */
    public function handle(): void
    {
        if ($this->batch()->cancelled()) {
            // Determine if the batch has been cancelled...

            return;
        }

        // Import a portion of the CSV file...
    }
}

分发批次

要分发一批任务,应使用 batch 方法,该方法由 Bus 门面提供。当然,批处理主要是在配合完成回调时发挥作用。因此,你可以使用 then、catch 和 finally 方法,为批次定义完成回调。这些回调在被调用时,都会收到一个 Illuminate\Bus\Batch 实例。

运行多个队列工作进程时,批次中的任务会被并行处理。因此,任务完成的顺序可能与加入批次的顺序不同。如果需要按顺序运行一系列任务,请参阅任务链与批次文档。

在这个示例中,假设我们将一批任务加入队列,每个任务负责处理 CSV 文件中的指定数量行:


use App\Jobs\ImportCsv;
use Illuminate\Bus\Batch;
use Illuminate\Support\Facades\Bus;
use Throwable;

$batch = Bus::batch([
    new ImportCsv(1, 100),
    new ImportCsv(101, 200),
    new ImportCsv(201, 300),
    new ImportCsv(301, 400),
    new ImportCsv(401, 500),
])->before(function (Batch $batch) {
    // The batch has been created but no jobs have been added...
})->progress(function (Batch $batch) {
    // A single job has completed successfully...
})->then(function (Batch $batch) {
    // All jobs completed successfully...
})->catch(function (Batch $batch, Throwable $e) {
    // Batch job failure detected...
})->finally(function (Batch $batch) {
    // The batch has finished executing...
})->dispatch();

return $batch->id;

批次 ID 可通过 $batch->id 属性取得。在批次分发后,可以使用该 ID 查询 Laravel 命令总线,获取批次信息。

由于批次回调会被序列化,并在稍后由 Laravel 队列执行,因此不应在回调中使用 $this 变量。此外,由于批处理任务会被包裹在数据库事务中,不应在任务内执行会触发隐式提交的数据库语句。

为批次命名

如果为批次命名,Laravel Horizon 和 Laravel Telescope 等工具可能会提供更易于理解的批次调试信息。要为批次指定任意名称,可以在定义批次时调用 name 方法:


$batch = Bus::batch([
    // ...
])->then(function (Batch $batch) {
    // All jobs completed successfully...
})->name('Import CSV')->dispatch();

批次的连接与队列

如果想指定批处理任务所使用的连接和队列,可以使用 onConnection 和 onQueue 方法。所有批处理任务都必须在同一个连接和队列中执行:


$batch = Bus::batch([
    // ...
])->then(function (Batch $batch) {
    // All jobs completed successfully...
})->onConnection('redis')->onQueue('imports')->dispatch();

任务链与批次

你可以将链式任务放入数组,在一个批次中定义一组任务链。例如,可以并行执行两条任务链,并在两条链都处理完毕后执行回调:


use App\Jobs\ReleasePodcast;
use App\Jobs\SendPodcastReleaseNotification;
use Illuminate\Bus\Batch;
use Illuminate\Support\Facades\Bus;

Bus::batch([
    [
        new ReleasePodcast(1),
        new SendPodcastReleaseNotification(1),
    ],
    [
        new ReleasePodcast(2),
        new SendPodcastReleaseNotification(2),
    ],
])->then(function (Batch $batch) {
    // All jobs completed successfully...
})->dispatch();

反过来,你也可以在任务链中定义批次,从而在链内运行任务批次。例如,可以先运行一批任务发布多个播客,然后再运行另一批任务发送发布通知:


use App\Jobs\FlushPodcastCache;
use App\Jobs\ReleasePodcast;
use App\Jobs\SendPodcastReleaseNotification;
use Illuminate\Support\Facades\Bus;

Bus::chain([
    new FlushPodcastCache,
    Bus::batch([
        new ReleasePodcast(1),
        new ReleasePodcast(2),
    ]),
    Bus::batch([
        new SendPodcastReleaseNotification(1),
        new SendPodcastReleaseNotification(2),
    ]),
])->dispatch();

向批次添加任务

有时,在某个批处理任务内部向其所属批次添加更多任务会很有用。当需要批处理数千个任务,而在一次 Web 请求期间分发它们可能耗时过长时,可以使用这种模式。因此,你可以先分发一批“加载器”任务,由这些任务向批次中填充更多任务:


$batch = Bus::batch([
    new LoadImportBatch,
    new LoadImportBatch,
    new LoadImportBatch,
])->then(function (Batch $batch) {
    // All jobs completed successfully...
})->name('Import Contacts')->dispatch();

在这个示例中,我们使用 LoadImportBatch 任务,向批次中填充更多任务。为此,可以调用批次实例的 add 方法;该实例可通过任务的 batch 方法取得:


use App\Jobs\ImportContacts;
use Illuminate\Support\Collection;

/**
 * Execute the job.
 */
public function handle(): void
{
    if ($this->batch()->cancelled()) {
        return;
    }

    $this->batch()->add(Collection::times(1000, function () {
        return new ImportContacts;
    }));
}

只有属于同一批次的任务,才能在其内部向该批次添加任务。

查看批次信息

传给批次完成回调的 Illuminate\Bus\Batch 实例提供了多种属性和方法,帮助你与指定任务批次交互,并查看其信息:


// The UUID of the batch...
$batch->id;

// The name of the batch (if applicable)...
$batch->name;

// The number of jobs assigned to the batch...
$batch->totalJobs;

// The number of jobs that have not been processed by the queue...
$batch->pendingJobs;

// The number of jobs that have failed...
$batch->failedJobs;

// The number of jobs that have been processed thus far...
$batch->processedJobs();

// The completion percentage of the batch (0-100)...
$batch->progress();

// Indicates if the batch has finished executing...
$batch->finished();

// Cancel the execution of the batch...
$batch->cancel();

// Indicates if the batch has been cancelled...
$batch->cancelled();

从路由返回批次

所有 Illuminate\Bus\Batch 实例都可序列化为 JSON。这意味着你可以直接从应用的某个路由返回这些实例,取得包含批次信息及其完成进度的 JSON 数据。这样就能方便地在应用界面中显示批次完成进度。

要根据 ID 获取批次,可以使用 Bus 门面的 findBatch 方法:


use Illuminate\Support\Facades\Bus;
use Illuminate\Support\Facades\Route;

Route::get('/batch/{batchId}', function (string $batchId) {
    return Bus::findBatch($batchId);
});

取消批次

有时,你可能需要取消某个批次的执行。可以通过调用 cancel 方法实现,该方法属于 Illuminate\Bus\Batch 实例:


/**
 * Execute the job.
 */
public function handle(): void
{
    if ($this->user->exceedsImportLimit()) {
        $this->batch()->cancel();

        return;
    }

    if ($this->batch()->cancelled()) {
        return;
    }
}

你可能已从前面的示例中注意到,批处理任务通常应在继续执行前,判断其所属批次是否已取消。不过,为方便起见,也可以直接为任务指定 SkipIfBatchCancelled 中间件。正如其名称所示,如果任务所属批次已取消,该中间件会指示 Laravel 不再处理该任务:


use Illuminate\Queue\Middleware\SkipIfBatchCancelled;

/**
 * Get the middleware the job should pass through.
 */
public function middleware(): array
{
    return [new SkipIfBatchCancelled];
}

批次失败

批处理任务失败时,会调用 catch 回调(如果已指定)。该回调只会在批次中第一个任务失败时调用。

允许失败

批次中的某个任务失败时,Laravel 会自动将该批次标记为“已取消”。如果需要,可以禁用这一行为,使单个任务失败不会自动将批次标记为已取消。为此,可以在分发批次时调用 allowFailures 方法:


$batch = Bus::batch([
    // ...
])->then(function (Batch $batch) {
    // All jobs completed successfully...
})->allowFailures()->dispatch();

你还可以选择向 allowFailures 方法传入闭包,每次有任务失败时都会执行该闭包:


$batch = Bus::batch([
    // ...
])->allowFailures(function (Batch $batch, $exception) {
    // Handle individual job failures...
})->dispatch();

重试批次中的失败任务

为方便起见,Laravel 提供了 queue:retry-batch Artisan 命令,让你可以轻松地重试某个批次中的所有失败任务。该命令接受要重试失败任务的批次 UUID:


php artisan queue:retry-batch 32dbc76c-4f82-4749-b610-a639fe0099b5

清理批次

如果不进行清理,job_batches 表可能很快就会积累大量记录。为缓解这个问题,应调度 queue:prune-batches Artisan 命令,使其每天运行:


use Illuminate\Support\Facades\Schedule;

Schedule::command('queue:prune-batches')->daily();

默认情况下,会清理所有已完成超过 24 小时的批次。调用该命令时,可以使用 hours 选项决定保留批次数据的时长。例如,以下命令会删除所有已完成超过 48 小时的批次:


use Illuminate\Support\Facades\Schedule;

Schedule::command('queue:prune-batches --hours=48')->daily();

有时,job_batches 表可能积累一些始终没有成功完成的批次记录,例如某个任务失败后,一直未能成功重试的批次。你可以让 queue:prune-batches 命令通过 unfinished 选项,清理这些未完成批次的记录:


use Illuminate\Support\Facades\Schedule;

Schedule::command('queue:prune-batches --hours=48 --unfinished=72')->daily();

同样,job_batches 表也可能积累已取消批次的记录。你可以让 queue:prune-batches 命令通过 cancelled 选项,清理这些已取消批次的记录:


use Illuminate\Support\Facades\Schedule;

Schedule::command('queue:prune-batches --hours=48 --cancelled=72')->daily();

在 DynamoDB 中存储批次

Laravel 还支持将批次元信息存储在 DynamoDB 中,而不是关系数据库中。不过,你需要手动创建一个 DynamoDB 表,用来保存所有批次记录。

通常,这个表应命名为 job_batches,但实际表名应以 queue.batching.table 配置值为准,该值位于应用的 queue 配置文件中。

DynamoDB 批次表配置

job_batches 表应具有一个名为 application 的字符串类型主分区键,以及一个名为 id 的字符串类型主排序键。键中的 application 部分保存应用名称,该名称由 name 配置值定义,此配置值位于应用的 app 配置文件中。由于应用名称是 DynamoDB 表键的一部分,因此你可以使用同一张表存储多个 Laravel 应用的任务批次。

此外,还可以为表定义 ttl 属性,以便利用自动批次清理功能。

DynamoDB 配置

接下来,安装 AWS SDK,使 Laravel 应用能够与 Amazon DynamoDB 通信:


composer require aws/aws-sdk-php

然后,将 queue.batching.driver 配置选项的值设为 dynamodb。此外,还应定义 key、secret 和 region 配置选项,并将它们放在 batching 配置数组中。这些选项用于向 AWS 进行身份验证。使用 dynamodb 驱动时,无须设置 queue.batching.database 配置选项:


'batching' => [
    'driver' => env('QUEUE_BATCHING_DRIVER', 'dynamodb'),
    'key' => env('AWS_ACCESS_KEY_ID'),
    'secret' => env('AWS_SECRET_ACCESS_KEY'),
    'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),
    'table' => 'job_batches',
],

清理 DynamoDB 中的批次

使用 DynamoDB 存储任务批次信息时,通常用于清理关系数据库中批次的命令不会生效。你可以改用 DynamoDB 原生的 TTL 功能,自动移除旧批次记录。

如果为 DynamoDB 表定义了 ttl 属性,就可以定义配置参数,指示 Laravel 如何清理批次记录。queue.batching.ttl_attribute 配置值定义保存 TTL 的属性名称,而 queue.batching.ttl 配置值定义从记录最后一次更新算起,需要经过多少秒,才能从 DynamoDB 表中移除该批次记录:


'batching' => [
    'driver' => env('QUEUE_BATCHING_DRIVER', 'dynamodb'),
    'key' => env('AWS_ACCESS_KEY_ID'),
    'secret' => env('AWS_SECRET_ACCESS_KEY'),
    'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),
    'table' => 'job_batches',
    'ttl_attribute' => 'ttl',
    'ttl' => 60 * 60 * 24 * 7, // 7 days...
],

将闭包加入队列

除了将任务类分发到队列,你也可以分发闭包。这非常适合需要在当前请求周期之外执行的简短、简单任务。将闭包分发到队列时,其代码内容会通过密码学方式签名,以防止传输过程中被修改:


use App\Models\Podcast;

$podcast = Podcast::find(1);

dispatch(function () use ($podcast) {
    $podcast->publish();
});

要为入队闭包指定名称,以便队列报告仪表盘使用,并在 queue:work 命令的输出中显示,可以使用 name 方法:


dispatch(function () {
    // ...
})->name('Publish Podcast');

使用 catch 方法,可以提供一个闭包:如果入队闭包用尽队列配置的重试次数后仍未成功完成,就会执行该闭包:


use Throwable;

dispatch(function () use ($podcast) {
    $podcast->publish();
})->catch(function (Throwable $e) {
    // This job has failed...
});

由于 catch 回调会被序列化,并由 Laravel 队列在稍后执行,因此不应将 $this 变量用于 catch 回调中。

运行队列工作进程

queue:work 命令

Laravel 提供了一个 Artisan 命令,用于启动队列工作进程,并处理推送到队列中的新任务。可以使用 queue:work Artisan 命令运行工作进程。请注意,queue:work 命令启动后会持续运行,直到手动停止或关闭终端:


php artisan queue:work

要让 queue:work 进程长期在后台运行,应使用 Supervisor 等进程监控工具,确保队列工作进程不会停止运行。

如果希望命令输出包含所处理任务的 ID、连接名称和队列名称,可以添加 -v 标志来调用 queue:work 命令:


php artisan queue:work -v

请记住,队列工作进程是长期运行的进程,会将已启动的应用状态保存在内存中。因此,启动之后,它们无法感知代码库的变更。所以,在部署过程中,一定要重启队列工作进程。此外,应用创建或修改的任何静态状态都不会在任务之间自动重置。

也可以运行 queue:listen 命令。使用 queue:listen 命令时,如果要重新加载更新后的代码或重置应用状态,无需手动重启工作进程;不过,该命令的效率明显低于 queue:work 命令:


php artisan queue:listen

运行多个队列工作进程

要为一个队列分配多个工作进程并并发处理任务,只需启动多个 queue:work 进程。可以在本地通过终端的多个标签页完成,也可以在生产环境中通过进程管理器的配置设置实现。使用 Supervisor 时,可以使用 numprocs 配置值。

指定连接和队列

还可以指定工作进程应使用的队列连接。传递给 work 命令的连接名称应对应 config/queue.php 配置文件中定义的某个连接:


php artisan queue:work redis

默认情况下,queue:work 命令只处理指定连接上的默认队列中的任务。不过,还可以进一步定制队列工作进程,使其仅处理某个连接上的特定队列。例如,如果所有邮件都由 emails 队列处理,且该队列位于 redis 队列连接上,可以执行以下命令,启动一个仅处理该队列的工作进程:


php artisan queue:work redis --queue=emails

处理指定数量的任务

可以使用 --once 选项,让工作进程仅处理队列中的一个任务:


php artisan queue:work --once

可以使用 --max-jobs 选项,让工作进程处理指定数量的任务后退出。该选项与 Supervisor 配合使用时很有帮助,可以让工作进程处理指定数量的任务后自动重启,释放运行期间可能积累的内存:


php artisan queue:work --max-jobs=1000

处理所有排队任务后退出

可以使用 --stop-when-empty 选项,让工作进程处理所有任务后正常退出。如果在 Docker 容器中处理 Laravel 队列,并希望队列清空后关闭容器,该选项会很有用:


php artisan queue:work --stop-when-empty

处理任务指定秒数后退出

可以使用 --max-time 选项,让工作进程处理任务指定秒数后退出。该选项与 Supervisor 配合使用时很有帮助,可以让工作进程处理任务达到指定时长后自动重启,释放运行期间可能积累的内存:


# Process jobs for one hour and then exit...
php artisan queue:work --max-time=3600

工作进程的休眠时长

当队列中有可处理的任务时,工作进程会持续处理任务,任务之间没有延迟。不过,sleep 选项决定了没有可处理任务时,工作进程将“休眠”多少秒。当然,休眠期间,工作进程不会处理任何新任务:


php artisan queue:work --sleep=3

维护模式与队列

应用处于维护模式时,不会处理任何排队任务。应用退出维护模式后,任务会恢复正常处理。

要强制队列工作进程在启用维护模式时仍然处理任务,可以使用 --force 选项:


php artisan queue:work --force

资源方面的注意事项

作为守护进程运行的队列工作进程不会在处理每个任务前“重新启动”框架。因此,每个任务完成后,都应释放占用较大的资源。例如,使用 图像处理功能和 GD 库时,应在处理完图像后使用 imagedestroy 释放内存。

队列优先级

有时可能需要为队列的处理设置优先级。例如,在 config/queue.php 配置文件中,可以将默认 queue 设置为 redis 连接上的 low 队列。不过,偶尔也可能需要像下面这样,将任务推送到优先级为 high 的队列:


dispatch((new Job)->onQueue('high'));

要启动一个工作进程,确保处理完所有 high 队列中的任务后,才开始处理 low 队列中的任务,可以向 work 命令传递以逗号分隔的队列名称列表:


php artisan queue:work --queue=high,low

队列工作进程与部署

由于队列工作进程是长期运行的进程,不重启就无法感知代码的变更。因此,部署使用队列工作进程的应用时,最简单的方式是在部署过程中重启这些工作进程。可以执行 queue:restart 命令,让所有工作进程正常重启:


php artisan queue:restart

该命令会通知所有队列工作进程,在完成当前任务后正常退出,从而避免丢失现有任务。由于执行 queue:restart 命令后,队列工作进程会退出,因此应运行 Supervisor 等进程管理器,以自动重启队列工作进程。

队列使用缓存存储重启信号,因此在使用此功能前,应确认应用已正确配置缓存驱动。

响应工作进程信号

队列工作进程在处理任务时收到 SIGQUIT、SIGTERM 或 SIGINT 等终止信号,会先完成当前任务再退出。不过,在服务器或容器编排系统停止进程之前,任务可能需要先响应该信号。例如,长时间运行的导入任务可能需要停止拉取新记录,并保存当前进度。

要在任务内部响应工作进程信号,请实现 Illuminate\Contracts\Queue\Interruptible 接口,并在任务中定义 interrupted 方法。工作进程收到的信号编号会传递给 interrupted 方法:


<?php

namespace App\Jobs;

use App\Models\Import;
use Illuminate\Contracts\Queue\Interruptible;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;

class ImportProducts implements ShouldQueue, Interruptible
{
    use Queueable;

    protected bool $shouldStop = false;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public Import $import,
    ) {}

    /**
     * Execute the job.
     */
    public function handle(): void
    {
        foreach ($this->import->pendingRows() as $row) {
            if ($this->shouldStop) {
                break;
            }

            // Import the product row...
        }

        $this->import->saveProgress();
    }

    /**
     * Handle a signal received by the queue worker.
     */
    public function interrupted(int $signal): void
    {
        $this->shouldStop = true;
    }
}

只有在任务正在运行期间,工作进程收到进程信号时,才会调用 interrupted 方法。它不能替代超时机制或任务的 failed 方法。

任务过期与超时

任务过期

在 config/queue.php 配置文件中,每个队列连接都定义了 retry_after 选项。该选项指定队列连接在重试正在处理的任务之前,应等待多少秒。例如,如果将 retry_after 的值设为 90,任务在处理了 90 秒后仍未被释放或删除,就会被重新释放回队列。通常,应将 retry_after 的值设为任务正常完成处理所需的最长合理秒数。

唯一不包含 retry_after 值的队列连接是 Amazon SQS。SQS 会根据在 AWS 控制台中管理的默认可见性超时重试任务。

工作进程超时

queue:work Artisan 命令提供了 --timeout 选项。默认情况下,--timeout 的值为 60 秒。如果任务的处理时间超过超时值指定的秒数,处理该任务的工作进程会以错误状态退出。通常,工作进程会由服务器上配置的进程管理器自动重启:


php artisan queue:work --timeout=60

retry_after 配置选项和 --timeout 命令行选项虽然不同,但会配合工作,确保任务不会丢失,并且仅成功处理一次。

--timeout 的值应始终比 retry_after 配置值至少短几秒。这样可以确保处理卡住任务的工作进程,总是在该任务被重试之前终止。如果 --timeout 选项的值大于 retry_after 配置值,任务可能会被处理两次。

暂停与恢复队列工作进程

有时可能需要暂时阻止队列工作进程处理新任务,而不完全停止该进程。例如,系统维护期间可能需要暂停任务处理。Laravel 提供了 queue:pause 和 queue:continue Artisan 命令,用于暂停和恢复队列工作进程。

要暂停某个特定队列,请提供队列连接名称和队列名称:


php artisan queue:pause database:default

在本例中,database 是队列连接名称,default 是队列名称。队列暂停后,正在处理该队列任务的工作进程会继续完成当前任务,但在队列恢复之前,不会领取任何新任务。

要暂停所有连接上所有队列的任务处理,请使用 --all 选项:


php artisan queue:pause --all

要恢复处理已暂停队列中的任务,请使用 queue:continue 命令:


php artisan queue:continue database:default

要恢复所有连接上所有队列的任务处理,请将 --all 选项与 queue:resume 命令一起使用:


php artisan queue:resume --all

队列恢复后,工作进程会立即开始处理该队列中的新任务。恢复所有队列不会恢复那些单独暂停的队列。请注意,暂停队列并不会停止工作进程本身,只会阻止它处理指定队列中的新任务。

工作进程的重启与暂停信号

默认情况下,队列工作进程在每轮任务处理时,都会轮询缓存驱动,检查重启和暂停信号。虽然这种轮询对于响应 queue:restart 和 queue:pause 命令必不可少,但它确实会带来少量性能开销。

如果需要优化性能,且不需要这些中断功能,可以调用 withoutInterruptionPolling 方法,在全局禁用轮询,该方法由 Queue 门面提供。通常应在 boot 方法中执行此操作,该方法位于 AppServiceProvider 中:


use Illuminate\Support\Facades\Queue;

/**
 * Bootstrap any application services.
 */
public function boot(): void
{
    Queue::withoutInterruptionPolling();
}

也可以通过设置静态属性 $restartable 或 $pausable,分别禁用重启轮询或暂停轮询。这些属性位于 Illuminate\Queue\Worker 类中:


use Illuminate\Queue\Worker;

/**
 * Bootstrap any application services.
 */
public function boot(): void
{
    Worker::$restartable = false;
    Worker::$pausable = false;
}

禁用中断轮询后,工作进程将不再响应 queue:restart 或 queue:pause 命令,具体取决于禁用了哪些功能。

Supervisor 配置

在生产环境中,需要一种方式来保持 queue:work 进程持续运行。queue:work 进程可能由于多种原因停止运行,例如工作进程超时,或者执行了 queue:restart 命令。

因此,需要配置一个进程监控工具,在 queue:work 进程退出时检测到这一情况,并自动重启它们。此外,进程监控工具还可以指定希望同时运行多少个 queue:work 进程。Supervisor 是 Linux 环境中常用的进程监控工具,下面将介绍如何配置它。

安装 Supervisor

Supervisor 是 Linux 操作系统的进程监控工具,能够在 queue:work 进程失败时自动重启它们。要在 Ubuntu 上安装 Supervisor,可以使用以下命令:


sudo apt-get install supervisor

如果自行配置和管理 Supervisor 让你觉得负担较重,可以考虑使用 Laravel Cloud,它提供了运行 Laravel 队列工作进程的全托管平台。

配置 Supervisor

Supervisor 配置文件通常存储在 /etc/supervisor/conf.d 目录中。可以在该目录下创建任意数量的配置文件,告诉 Supervisor 应如何监控进程。例如,创建一个 laravel-worker.conf 文件,用于启动并监控 queue:work 进程:


[program:laravel-worker]
process_name=%(program_name)s_%(process_num)02d
command=php /home/forge/app.com/artisan queue:work --sleep=3 --tries=3 --max-time=3600
autostart=true
autorestart=true
stopasgroup=true
killasgroup=true
user=forge
numprocs=8
redirect_stderr=true
stdout_logfile=/home/forge/app.com/worker.log
stopwaitsecs=3600

在本例中,numprocs 指令会让 Supervisor 运行并监控八个 queue:work 进程,在这些进程失败时自动重启它们。应修改配置中的 command 指令,以使用所需的队列连接和工作进程选项。

应确保 stopwaitsecs 的值大于运行时间最长的任务所耗费的秒数。否则,Supervisor 可能会在任务处理完成前将其终止。

启动 Supervisor

创建配置文件后,可以使用以下命令更新 Supervisor 配置,并启动进程:


sudo supervisorctl reread

sudo supervisorctl update

sudo supervisorctl start "laravel-worker:*"

有关 Supervisor 的更多信息,请参阅 Supervisor 文档。

处理失败任务

排队任务有时会失败。无需担心,事情并不总会按计划进行!Laravel 提供了一种便捷方式来指定任务的最大尝试次数。异步任务超过该次数后,会被记录到 failed_jobs 数据库表中。同步派发的任务如果失败,不会存入此表,其异常会立即由应用处理。

新的 Laravel 应用通常已经包含用于创建 failed_jobs 表的迁移。不过,如果应用中没有该表的迁移,可以使用 make:queue-failed-table 命令创建迁移:


php artisan make:queue-failed-table

php artisan migrate

运行队列工作进程时,可以使用 --tries 开关指定任务的最大尝试次数,该开关由 queue:work 命令提供。如果未为 --tries 选项指定值,任务只会尝试一次,或者按任务类的 Tries 特性所指定的次数尝试:


php artisan queue:work redis --tries=3

使用 --backoff 选项,可以指定 Laravel 在重试发生异常的任务之前,应等待多少秒。默认情况下,任务会立即被释放回队列,以便再次尝试:


php artisan queue:work redis --tries=3 --backoff=3

如果希望针对每个任务配置 Laravel 在重试发生异常的任务之前应等待的秒数,可以在任务类上使用 Backoff 特性:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Backoff;

#[Backoff(3)]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

如果需要更复杂的逻辑来确定任务的重试等待时间,可以在任务类上定义 backoff 方法:


/**
 * Calculate the number of seconds to wait before retrying the job.
 */
public function backoff(): int
{
    return 3;
}

通过定义一组重试等待时间,可以轻松配置“指数退避”。在本例中,第一次重试的延迟为 1 秒,第二次为 5 秒,第三次为 10 秒;如果还有剩余尝试次数,之后每次重试的延迟也都是 10 秒:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Backoff;

#[Backoff([1, 5, 10])]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

任务失败后的清理

某个任务失败时,可能需要向用户发送提醒,或者撤销该任务已部分完成的操作。为此,可以在任务类上定义 failed 方法。导致任务失败的 Throwable 实例会传递给 failed 方法:


<?php

namespace App\Jobs;

use App\Models\Podcast;
use App\Services\AudioProcessor;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Queue\Queueable;
use Throwable;

class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * Create a new job instance.
     */
    public function __construct(
        public Podcast $podcast,
    ) {}

    /**
     * Execute the job.
     */
    public function handle(AudioProcessor $processor): void
    {
        // Process uploaded podcast...
    }

    /**
     * Handle a job failure.
     */
    public function failed(?Throwable $exception): void
    {
        // Send user notification of failure, etc...
    }
}

调用 failed 方法前,会重新实例化一个任务对象。因此,在 handle 方法中对类属性所做的任何修改都会丢失。

失败的任务不一定是遇到了未处理异常的任务。当任务用尽所有允许的尝试次数时,也可能被视为失败。尝试次数可能以多种方式消耗:

  • 任务超时。
  • 任务在执行期间遇到未处理的异常。
  • 任务被手动释放回队列,或由中间件释放回队列。

如果最后一次尝试因为任务执行期间抛出的异常而失败,该异常会传递给任务的 failed 方法。不过,如果任务因为达到允许的最大尝试次数而失败,$exception 会是 Illuminate\Queue\MaxAttemptsExceededException 的实例。同样,如果任务因为超过配置的超时时间而失败,$exception 会是 Illuminate\Queue\TimeoutExceededException 的实例。

重试失败任务

要查看已存入 failed_jobs 数据库表的所有失败任务,可以使用 queue:failed Artisan 命令:


php artisan queue:failed

queue:failed 命令会列出任务 ID、连接、队列、失败时间及其他任务信息。可以使用任务 ID 来重试失败任务。例如,要重试 ID 为 ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece 的失败任务,请执行以下命令:


php artisan queue:retry ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece

如果需要,可以向该命令传递多个 ID:


php artisan queue:retry ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece 91401d2c-0784-4f43-824c-34f94a33c24d

也可以重试某个特定队列中的所有失败任务:


php artisan queue:retry --queue=name

要重试所有失败任务,请执行 queue:retry 命令,并将 all 作为 ID 传入:


php artisan queue:retry all

如果要删除某个失败任务,可以使用 queue:forget 命令:


php artisan queue:forget 91401d2c-0784-4f43-824c-34f94a33c24d

使用 Horizon 时,应使用 horizon:forget 命令删除失败任务,而不是使用 queue:forget 命令。

要删除 failed_jobs 表中的所有失败任务,可以使用 queue:flush 命令:


php artisan queue:flush

queue:flush 命令会从队列中删除所有失败任务记录,无论这些任务已失败多久。可以使用 --hours 选项,仅删除在指定小时数之前或更早失败的任务:


php artisan queue:flush --hours=48

忽略缺失的模型

向任务注入 Eloquent 模型时,模型会在任务放入队列前自动序列化,并在任务处理时重新从数据库获取。不过,如果模型在任务等待工作进程处理期间被删除,任务可能会因 ModelNotFoundException 而失败。

为方便起见,可以在任务类上使用 DeleteWhenMissingModels 特性,让模型缺失的任务自动删除。存在该特性时,Laravel 会直接丢弃任务,不抛出异常:


<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\DeleteWhenMissingModels;

#[DeleteWhenMissingModels]
class ProcessPodcast implements ShouldQueue
{
    // ...
}

清理过期失败任务

要清理应用 failed_jobs 表中的记录,可以调用 queue:prune-failed Artisan 命令:


php artisan queue:prune-failed

默认情况下,会清理所有超过 24 小时的失败任务记录。如果为该命令提供 --hours 选项,则只保留最近 N 小时内插入的失败任务记录。例如,以下命令会删除所有在 48 小时以前插入的失败任务记录:


php artisan queue:prune-failed --hours=48

将失败任务存储在 DynamoDB 中

Laravel 还支持将失败任务记录存储在 DynamoDB 中,而不是关系型数据库表中。不过,必须手动创建一个 DynamoDB 表来存储所有失败任务记录。通常,该表应命名为 failed_jobs,但实际表名应根据 queue.failed.table 配置值来确定,该配置值位于应用的 queue 配置文件中。

failed_jobs 表应包含一个名为 application 的字符串类型主分区键,以及一个名为 uuid 的字符串类型主排序键。键的 application 部分将包含应用名称,该名称由 name 配置值定义,此配置值位于应用的 app 配置文件中。由于应用名称是 DynamoDB 表键的一部分,因此可以使用同一个表存储多个 Laravel 应用的失败任务。

此外,请确保安装 AWS SDK,以便 Laravel 应用能够与 Amazon DynamoDB 通信:


composer require aws/aws-sdk-php

接下来,将 queue.failed.driver 配置选项的值设为 dynamodb。此外,还应在失败任务配置数组中定义 key、secret 和 region 配置选项。这些选项将用于 AWS 身份验证。使用 dynamodb 驱动时,不需要 queue.failed.database 配置选项:


'failed' => [
    'driver' => env('QUEUE_FAILED_DRIVER', 'dynamodb'),
    'key' => env('AWS_ACCESS_KEY_ID'),
    'secret' => env('AWS_SECRET_ACCESS_KEY'),
    'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),
    'table' => 'failed_jobs',
],

禁用失败任务存储

可以将 queue.failed.driver 配置选项的值设为 null,让 Laravel 丢弃失败任务而不存储它们。通常可以通过 QUEUE_FAILED_DRIVER 环境变量完成此设置:


QUEUE_FAILED_DRIVER=null

失败任务事件

如果希望注册一个在任务失败时调用的事件监听器,可以使用 Queue 门面的 failing 方法。例如,可以在 boot 方法中为该事件添加一个闭包,该方法位于 Laravel 自带的 AppServiceProvider 中:


<?php

namespace App\Providers;

use Illuminate\Support\Facades\Queue;
use Illuminate\Support\ServiceProvider;
use Illuminate\Queue\Events\JobFailed;

class AppServiceProvider extends ServiceProvider
{
    /**
     * Register any application services.
     */
    public function register(): void
    {
        // ...
    }

    /**
     * Bootstrap any application services.
     */
    public function boot(): void
    {
        Queue::failing(function (JobFailed $event) {
            // $event->connectionName
            // $event->job
            // $event->exception
        });
    }
}

清空队列中的任务

使用 Horizon 时,应使用 horizon:clear 命令清空队列中的任务,而不是使用 queue:clear 命令。

如果希望删除默认连接的默认队列中的所有任务,可以使用 queue:clear Artisan 命令:


php artisan queue:clear

也可以提供 connection 参数和 queue 选项,删除特定连接和队列中的任务:


php artisan queue:clear redis --queue=emails

只有 SQS、Redis 和数据库队列驱动支持清空队列中的任务。此外,SQS 的消息删除过程最多需要 60 秒,因此在清空队列之后的 60 秒内发送到 SQS 队列的任务,也可能被删除。

监控队列

如果队列突然收到大量任务,可能会不堪重负,导致任务等待较长时间才能完成。如果需要,Laravel 可以在队列任务数量超过指定阈值时发出提醒。

首先,应将 queue:monitor 命令安排为每分钟运行一次。该命令接受要监控的队列名称,以及所需的任务数量阈值:


php artisan queue:monitor redis:default,redis:deployments --max=100

仅调度该命令,还不足以触发队列负载过高的通知。当命令发现某个队列的任务数量超过阈值时,会派发 Illuminate\Queue\Events\QueueBusy 事件。可以在应用的 AppServiceProvider 中监听该事件,以便向自己或开发团队发送通知:


use App\Notifications\QueueHasLongWaitTime;
use Illuminate\Queue\Events\QueueBusy;
use Illuminate\Support\Facades\Event;
use Illuminate\Support\Facades\Notification;

/**
 * Bootstrap any application services.
 */
public function boot(): void
{
    Event::listen(function (QueueBusy $event) {
        Notification::route('mail', '[email protected]')
            ->notify(new QueueHasLongWaitTime(
                $event->connectionName,
                $event->queue,
                $event->size
            ));
    });
}

测试

测试派发任务的代码时,可能希望让 Laravel 不实际执行任务本身,因为任务的代码可以直接测试,并与派发任务的代码分开测试。当然,要测试任务本身,可以实例化一个任务对象,并在测试中直接调用 handle 方法。

可以使用 Queue 门面的 fake 方法,阻止排队任务实际推送到队列。调用 Queue 门面的 fake 方法后,就可以断言应用曾尝试将任务推送到队列:


<?php

use App\Jobs\AnotherJob;
use App\Jobs\ShipOrder;
use Illuminate\Support\Facades\Queue;

test('orders can be shipped', function () {
    Queue::fake();

    // Perform order shipping...

    // Assert that no jobs were pushed...
    Queue::assertNothingPushed();

    // Assert a job was pushed to a given queue...
    Queue::assertPushedOn('queue-name', ShipOrder::class);

    // Assert a job was pushed
    Queue::assertPushed(ShipOrder::class);

    // Assert a job was pushed exactly once...
    Queue::assertPushedOnce(ShipOrder::class);

    // Assert a job was pushed twice...
    Queue::assertPushedTimes(ShipOrder::class, 2);

    // Assert a job was not pushed...
    Queue::assertNotPushed(AnotherJob::class);

    // Assert that a closure was pushed to the queue...
    Queue::assertClosurePushed();

    // Assert that a closure was not pushed...
    Queue::assertClosureNotPushed();

    // Assert the total number of jobs that were pushed...
    Queue::assertCount(3);
});

<?php

namespace Tests\Feature;

use App\Jobs\AnotherJob;
use App\Jobs\ShipOrder;
use Illuminate\Support\Facades\Queue;
use Tests\TestCase;

class ExampleTest extends TestCase
{
    public function test_orders_can_be_shipped(): void
    {
        Queue::fake();

        // Perform order shipping...

        // Assert that no jobs were pushed...
        Queue::assertNothingPushed();

        // Assert a job was pushed to a given queue...
        Queue::assertPushedOn('queue-name', ShipOrder::class);

        // Assert a job was pushed
        Queue::assertPushed(ShipOrder::class);

        // Assert a job was pushed exactly once...
        Queue::assertPushedOnce(ShipOrder::class);

        // Assert a job was pushed twice...
        Queue::assertPushedTimes(ShipOrder::class, 2);

        // Assert a job was not pushed...
        Queue::assertNotPushed(AnotherJob::class);

        // Assert that a closure was pushed to the queue...
        Queue::assertClosurePushed();

        // Assert that a closure was not pushed...
        Queue::assertClosureNotPushed();

        // Assert the total number of jobs that were pushed...
        Queue::assertCount(3);
    }
}

可以向 assertPushed、assertNotPushed、assertClosurePushed 或 assertClosureNotPushed 方法传入闭包,断言有任务被推送,并且该任务满足给定的判定条件。只要至少有一个被推送的任务满足该条件,断言就会成功:


use Illuminate\Queue\CallQueuedClosure;

Queue::assertPushed(function (ShipOrder $job) use ($order) {
    return $job->order->id === $order->id;
});

Queue::assertClosurePushed(function (CallQueuedClosure $job) {
    return $job->name === 'validate-order';
});

仅模拟部分任务

如果只需模拟特定任务,同时让其他任务正常执行,可以将需要模拟的任务类名传递给 fake 方法:


test('orders can be shipped', function () {
    Queue::fake([
        ShipOrder::class,
    ]);

    // Perform order shipping...

    // Assert a job was pushed twice...
    Queue::assertPushedTimes(ShipOrder::class, 2);
});

public function test_orders_can_be_shipped(): void
{
    Queue::fake([
        ShipOrder::class,
    ]);

    // Perform order shipping...

    // Assert a job was pushed twice...
    Queue::assertPushedTimes(ShipOrder::class, 2);
}

可以使用 except 方法模拟除一组指定任务之外的所有任务:


Queue::fake()->except([
    ShipOrder::class,
]);

测试任务链

要测试任务链,需要使用 Bus 门面的模拟功能。Bus 门面的 assertChained 方法可用于断言某个任务链已被派发。assertChained 方法接受一个由链式任务组成的数组作为第一个参数:


use App\Jobs\RecordShipment;
use App\Jobs\ShipOrder;
use App\Jobs\UpdateInventory;
use Illuminate\Support\Facades\Bus;

Bus::fake();

// ...

Bus::assertChained([
    ShipOrder::class,
    RecordShipment::class,
    UpdateInventory::class
]);

如上例所示,链式任务数组可以是任务类名的数组。不过,也可以提供实际任务实例的数组。这样做时,Laravel 会确保任务实例与应用派发的链式任务属于相同的类,并且拥有相同的属性值:


Bus::assertChained([
    new ShipOrder,
    new RecordShipment,
    new UpdateInventory,
]);

可以使用 assertDispatchedWithoutChain 方法,断言任务被推送时没有附带任务链:


Bus::assertDispatchedWithoutChain(ShipOrder::class);

测试任务链的修改

如果链式任务向现有任务链的开头或末尾添加任务,可以使用该任务的 assertHasChain 方法,断言其剩余任务链符合预期:


$job = new ProcessPodcast;

$job->handle();

$job->assertHasChain([
    new TranscribePodcast,
    new OptimizePodcast,
    new ReleasePodcast,
]);

可以使用 assertDoesntHaveChain 方法,断言任务的剩余任务链为空:


$job->assertDoesntHaveChain();

测试任务链中的批次

如果任务链包含一个任务批次,可以在任务链断言中插入 Bus::chainedBatch 定义,以断言链中的批次符合预期:


use App\Jobs\ShipOrder;
use App\Jobs\UpdateInventory;
use Illuminate\Bus\PendingBatch;
use Illuminate\Support\Facades\Bus;

Bus::assertChained([
    new ShipOrder,
    Bus::chainedBatch(function (PendingBatch $batch) {
        return $batch->jobs->count() === 3;
    }),
    new UpdateInventory,
]);

测试任务批次

Bus 门面的 assertBatched 方法可用于断言某个任务批次已被派发。传递给 assertBatched 方法的闭包会收到一个 Illuminate\Bus\PendingBatch 实例,可用于检查批次中的任务:


use Illuminate\Bus\PendingBatch;
use Illuminate\Support\Facades\Bus;

Bus::fake();

// ...

Bus::assertBatched(function (PendingBatch $batch) {
    return $batch->name == 'Import CSV' &&
           $batch->jobs->count() === 10;
});

可以在待派发批次上使用 hasJobs 方法,验证批次中是否包含预期的任务。该方法接受一个由任务实例、类名或闭包组成的数组:


Bus::assertBatched(function (PendingBatch $batch) {
    return $batch->hasJobs([
        new ProcessCsvRow(row: 1),
        new ProcessCsvRow(row: 2),
        new ProcessCsvRow(row: 3),
    ]);
});

使用闭包时,闭包会收到任务实例。预期的任务类型将根据闭包的类型提示推断:


Bus::assertBatched(function (PendingBatch $batch) {
    return $batch->hasJobs([
        fn (ProcessCsvRow $job) => $job->row === 1,
        fn (ProcessCsvRow $job) => $job->row === 2,
        fn (ProcessCsvRow $job) => $job->row === 3,
    ]);
});

可以使用 assertBatchCount 方法,断言已派发指定数量的批次:


Bus::assertBatchCount(3);

可以使用 assertNothingBatched,断言未派发任何批次:


Bus::assertNothingBatched();

测试任务与批次的交互

此外,有时还需要测试单个任务与其所属批次的交互。例如,可能需要测试某个任务是否取消了其批次的后续处理。为此,需要通过 withFakeBatch 方法,为任务分配一个模拟批次。withFakeBatch 方法会返回一个包含任务实例和模拟批次的元组:


[$job, $batch] = (new ShipOrder)->withFakeBatch();

$job->handle();

$this->assertTrue($batch->cancelled());
$this->assertEmpty($batch->added);

测试任务与队列的交互

有时可能需要测试排队任务是否将自身释放回队列,或者测试任务是否删除了自身。可以实例化任务,并调用 withFakeQueueInteractions 方法来测试这些队列交互。

模拟任务的队列交互之后,就可以调用任务上的 handle 方法。调用任务后,可以使用多种断言方法来验证任务与队列的交互:


use App\Exceptions\CorruptedAudioException;
use App\Jobs\ProcessPodcast;

$job = (new ProcessPodcast)->withFakeQueueInteractions();

$job->handle();

$job->assertReleased(delay: 30);
$job->assertDeleted();
$job->assertNotDeleted();
$job->assertFailed();
$job->assertFailedWith(CorruptedAudioException::class);
$job->assertNotFailed();

任务事件

使用 before 和 after 方法,可以指定在排队任务处理前或处理后执行的回调;这些方法由 Queue 门面提供。这些回调适合用于记录额外日志,或为仪表盘累加统计数据。通常,应在 boot 方法中调用这些方法,该方法位于服务提供者中。例如,可以使用 Laravel 自带的 AppServiceProvider:


<?php

namespace App\Providers;

use Illuminate\Support\Facades\Queue;
use Illuminate\Support\ServiceProvider;
use Illuminate\Queue\Events\JobProcessed;
use Illuminate\Queue\Events\JobProcessing;

class AppServiceProvider extends ServiceProvider
{
    /**
     * Register any application services.
     */
    public function register(): void
    {
        // ...
    }

    /**
     * Bootstrap any application services.
     */
    public function boot(): void
    {
        Queue::before(function (JobProcessing $event) {
            // $event->connectionName
            // $event->job
            // $event->job->payload()
        });

        Queue::after(function (JobProcessed $event) {
            // $event->connectionName
            // $event->job
            // $event->job->payload()
        });
    }
}

使用 looping 方法,可以指定在工作进程尝试从队列获取任务之前执行的回调;该方法由 Queue 门面提供。例如,可以注册一个闭包,回滚先前失败任务遗留的、尚未结束的事务:


use Illuminate\Support\Facades\DB;
use Illuminate\Support\Facades\Queue;

Queue::looping(function () {
    while (DB::transactionLevel() > 0) {
        DB::rollBack();
    }
});

当队列工作进程无法从队列获取任务时,Laravel 还会派发 Illuminate\Queue\Events\WorkerIdle 事件:


use Illuminate\Queue\Events\WorkerIdle;
use Illuminate\Support\Facades\Event;

Event::listen(function (WorkerIdle $event) {
    // $event->connectionName
    // $event->queue
    // $event->workerOptions
});

来源:Taylor Otwell 与 Laravel 文档贡献者,Queues,Laravel 13.x 官方文档。本稿基于用户提供的冻结正文完整汉化;代码缩进从同一官方源页的复制文本恢复,200 个代码块逐字保留。示例未在本机 Laravel 环境运行,涉及驱动、扩展、进程管理器及版本的行为应按项目实际环境核对。

许可:Laravel 文档仓库的 MIT License,Copyright (c) Taylor Otwell。以下保留许可原文。

MIT 许可原文
The MIT License (MIT)

Copyright (c) Taylor Otwell

Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the “Software”), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:

The above copyright notice and this permission notice shall be included in
all copies or substantial portions of the Software.

THE SOFTWARE IS PROVIDED “AS IS”, WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
THE SOFTWARE.

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容