Skip to content
全部文档

队列

简介

在构建 Web 应用程序时,你可能需要执行一些任务,例如解析和存储上传的 CSV 文件,这些任务在典型的 Web 请求期间需要很长时间才能执行。值得庆幸的是,Laravel 允许你轻松创建可以在后台处理的排队任务。通过将时间密集型任务移至队列,你的应用程序可以以极快的速度响应 Web 请求,并为你的客户提供更好的用户体验。

Laravel 队列提供跨各种不同队列后端的统一排队 API,例如 Amazon SQSRedis,甚至是关系数据库。

Laravel 的队列配置选项存储在应用程序的 config/queue.php 配置文件中。在此文件中,你将找到框架中包含的每个队列驱动程序的连接配置,包括数据库、Amazon SQSRedisBeanstalkd 驱动程序,以及将立即执行作业的同步驱动程序(在开发或测试期间使用)。还包括 null 队列驱动程序,它会丢弃排队的作业。

INFO

Laravel Horizon 是一个漂亮的仪表板和配置系统,适用于 Redis 驱动的队列。查看完整的 Horizon documentation 了解更多信息。

连接与队列

在开始使用 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 队列 Worker允许你指定应按优先级处理哪些队列。例如,如果你将作业推送到 high 队列,你可以运行一个为它们提供更高处理优先级的工作线程:

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

驱动说明与前置条件

数据库

为了使用 database 队列驱动程序,你将需要一个数据库表来保存作业。通常,这包含在 Laravel 的默认 0001_01_01_000002_create_jobs_table.php 数据库迁移 中;但是,如果你的应用程序不包含此迁移,你可以使用 make:queue-table Artisan 命令来创建它:

shell
php artisan make:queue-table

php artisan migrate

Redis

为了使用 redis 队列驱动程序,你应该在 config/database.php 配置文件中配置 Redis 数据库连接。

WARNING

redis 队列驱动程序不支持 serializercompression Redis 选项。

Redis 集群

如果你的 Redis 队列连接使用 Redis Cluster,则你的队列名称必须包含 key hash tag。这是为了确保给定队列的所有 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,
],

WARNING

block_for 设置为 0 将导致队列 Worker无限期阻塞,直到有作业可用。这也将阻止处理诸如 SIGTERM 之类的信号,直到处理下一个作业。

其他驱动前置条件

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

  • Amazon SQS: `aws/aws-sdk-php ~3.0`
  • Beanstalkd: `pda/pheanstalk ~5.0`
  • Redis: `predis/predis ~2.0` or phpredis PHP extension
  • [MongoDB](https://www.mongodb.com/docs/drivers/php/laravel-mongodb/current/queues/): `mongodb/laravel-mongodb`

创建任务

生成任务类

默认情况下,应用程序的所有可排队任务都存储在 app/Jobs 目录中。如果 app/Jobs 目录不存在,则在运行 make:job Artisan 命令时将创建该目录:

shell
php artisan make:job ProcessPodcast

生成的类将实现 Illuminate\Contracts\Queue\ShouldQueue 接口,向 Laravel 指示作业应被推送到队列中以异步运行。

INFO

作业存根可以使用 stub publishing 进行定制。

类结构

作业类非常简单,通常只包含 handle 方法,该方法在队列处理作业时调用。首先,让我们看一个示例作业类。在此示例中,我们假设我们管理播客发布服务,并且需要在发布之前处理上传的播客文件:

php
<?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 model 直接传递到排队任务的构造函数中。由于作业正在使用 Queueable 特征,因此当作业处理时,Eloquent 模型及其加载的关系将被正常序列化和反序列化。

如果你的排队任务在其构造函数中接受 Eloquent 模型,则只有该模型的标识符才会序列化到队列中。当作业实际处理时,队列系统将自动从数据库中重新检索完整的模型实例及其加载的关系。这种模型序列化方法允许将更小的作业负载发送到队列驱动程序。

handle Method Dependency Injection

当队列处理作业时,会调用 handle 方法。请注意,我们可以键入提示对作业的 handle 方法的依赖关系。 Laravel 服务容器 自动注入这些依赖项。

如果你想完全控制容器如何将依赖项注入 handle 方法,你可以使用容器的 bindMethod 方法。 bindMethod 方法接受接收作业和容器的回调。在回调中,你可以随意调用 handle 方法。通常,你应该从 App\Providers\AppServiceProvider 服务提供者boot 方法中调用此方法:

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));
});

WARNING

二进制数据(例如原始图像内容)应先通过 base64_encode 函数传递,然后再传递到排队任务。否则,作业在放入队列时可能无法正确序列化为 JSON。

入队的关联

由于当作业排队时,所有加载的 Eloquent 模型关系也会被序列化,因此序列化的作业字符串有时会变得非常大。此外,当反序列化作业并从数据库重新检索模型关系时,将完整检索它们。反序列化作业时,将不会应用在作业排队过程中序列化模型之前应用的任何先前关系约束。因此,如果你希望使用给定关系的子集,则应该在排队任务中重新约束该关系。

或者,为了防止关系被序列化,你可以在设置属性值时调用模型上的 withoutRelations 方法。此方法将返回模型的实例,但不带其加载的关系:

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

如果你使用 PHP constructor property promotion 并希望指示 Eloquent 模型不应序列化其关系,则可以使用 WithoutRelations 属性:

use Illuminate\Queue\Attributes\WithoutRelations;
php
/**
 * Create a new job instance.
 */
public function __construct(
    #[WithoutRelations]
    public Podcast $podcast,
) {}

如果作业接收 Eloquent 模型的集合或数组而不是单个模型,则在反序列化和执行作业时,该集合中的模型将不会恢复其关系。这是为了防止处理大量模型的作业过度使用资源。

唯一任务

WARNING

唯一作业需要支持的缓存驱动。目前,memcachedredisdynamodbdatabasefilearray 缓存驱动支持原子锁。此外,唯一作业约束不适用于批处理中的作业。

有时,你可能希望确保在任何时间点队列中只有特定作业的一个实例。你可以通过在作业类上实现 ShouldBeUnique 接口来实现这一点。该接口不需要你在类上定义任何其他方法:

php
<?php

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

class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    ...
}

在上面的示例中,UpdateSearchIndex 作业是唯一的。因此,如果作业的另一个实例已在队列中并且尚未完成处理,则不会调度该作业。

在某些情况下,你可能想定义特定的「键」使任务唯一,或指定超时时间,超时后任务不再保持唯一。为此,可在任务类上定义 uniqueIduniqueFor 属性或方法:

php
<?php

use App\Models\Product;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Contracts\Queue\ShouldBeUnique;

class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    /**
     * The product instance.
     *
     * @var \App\Product
     */
    public $product;

    /**
     * The number of seconds after which the job's unique lock will be released.
     *
     * @var int
     */
    public $uniqueFor = 3600;

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

在上面的示例中,UpdateSearchIndex 作业通过产品 ID 是唯一的。因此,具有相同产品 ID 的作业的任何新调度都将被忽略,直到现有作业完成处理。此外,如果现有作业在一小时内未处理,则唯一锁将被释放,并且可以将具有相同唯一键的另一个作业分派到队列中。

WARNING

如果你的应用程序从多个 Web 服务器或容器分派作业,则应确保所有服务器都与同一中央缓存服务器通信,以便 Laravel 可以准确确定作业是否唯一。

在开始处理前保持任务唯一

默认情况下,在作业完成处理或所有重试尝试失败后,唯一作业将被「解锁」。但是,在某些情况下,你可能希望在处理作业之前立即解锁作业。为了实现这一点,你的工作应该实现 ShouldBeUniqueUntilProcessing 合约而不是 ShouldBeUnique 合约:

php
<?php

use App\Models\Product;
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');
    }
}

INFO

如果你只需要限制作业的并发处理,请改用 WithoutOverlapping 作业中间件。

加密任务

Laravel 允许你通过 encryption 确保作业数据的隐私性和完整性。首先,只需将 ShouldBeEncrypted 接口添加到作业类即可。一旦将此接口添加到类中,Laravel 将在将你的作业推送到队列之前自动对其进行加密:

php
<?php

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

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

任务中间件

作业中间件允许你围绕排队任务的执行包装自定义逻辑,从而减少作业本身的样板代码。例如,考虑以下 handle 方法,该方法利用 Laravel 的 Redis 速率限制功能,仅允许每五秒处理一个作业:

use Illuminate\Support\Facades\Redis;
php
/**
 * 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 方法中做限流,不如定义处理限流的作业中间件。Laravel 没有规定作业中间件的默认位置,你可以放在应用中的任意位置。本例中我们将其放在 app/Jobs/Middleware 目录:

php
<?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);
            });
    }
}

正如你所看到的,与 route middleware 一样,作业中间件接收正在处理的作业以及应调用以继续处理作业的回调。

创建作业中间件后,可通过从作业的 middleware 方法返回它们来附加到作业。由 make:job Artisan 命令生成的作业没有该方法,因此需要手动添加到作业类中:

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

INFO

作业中间件还可以分配给 queueable event listenersmailablesnotifications

速率限制

虽然我们刚刚演示了如何编写你自己的速率限制作业中间件,但 Laravel 实际上包含一个速率限制中间件,你可以利用它来限制作业的速率。与 route rate limiters 一样,作业速率限制器是使用 RateLimiter 门面的 for 方法定义的。

例如,你可能希望允许用户每小时备份一次数据,但对高级客户不施加此类限制。为此,你可以在 AppServiceProviderboot 方法中定义 RateLimiter

use Illuminate\Cache\RateLimiting\Limit;
use Illuminate\Support\Facades\RateLimiter;
php
/**
 * 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;
php
/**
 * Get the middleware the job should pass through.
 *
 * @return array<int, object>
 */
public function middleware(): array
{
    return [new RateLimited('backups')];
}

将受限流的任务释放回队列仍会增加该任务的 attempts 总数。你可能需要相应调整任务类上的 triesmaxExceptions 属性;也可以使用 retryUntil 方法 定义任务不再被尝试的截止时间。

如果你不希望某个作业在速率受限时重试,你可以使用 dontRelease 方法:

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

INFO

如果你使用的是 Redis,则可以使用 Illuminate\Queue\Middleware\RateLimitedWithRedis 中间件,该中间件针对 Redis 进行了微调,并且比基本限速中间件更高效:

防止任务重叠

Laravel 包含 Illuminate\Queue\Middleware\WithoutOverlapping 中间件,可让你防止基于任意键的作业重叠。当排队任务正在修改一次只能由一个作业修改的资源时,这会很有帮助。

例如,假设你有一个更新用户信用评分的排队任务,并且你希望防止同一用户 ID 的信用评分更新作业重叠。为此,你可以从作业的 middleware 方法返回 WithoutOverlapping 中间件:

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

任何相同类型的重叠作业将被释放回队列。你还可以指定再次尝试释放的作业之前必须经过的秒数:

php
/**
 * 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 方法:

php
/**
 * 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 锁:

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

WARNING

WithoutOverlapping 中间件需要支持 locks 的缓存驱动程序。目前,memcachedredisdynamodbdatabasefilearray 缓存驱动程序支持原子锁。

跨任务类共享锁键

默认情况下,WithoutOverlapping 中间件只会阻止同一类的重叠作业。因此,尽管两个不同的作业类别可能使用相同的锁定密钥,但不会阻止它们重叠。但是,你可以指示 Laravel 使用 shared 方法跨作业类应用密钥:

php
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 交互的排队任务开始抛出异常。要限制异常,你可以从作业的 middleware 方法返回 ThrottlesExceptions 中间件。通常,该中间件应与实现 time based attempts 的作业配对:

use DateTime;
use Illuminate\Queue\Middleware\ThrottlesExceptions;
php
/**
 * 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()->addMinutes(30);
}

中间件接受的第一个构造函数参数是作业在被限制之前可以抛出的异常数,而第二个构造函数参数是作业被限制后再次尝试之前应该经过的秒数。在上面的代码示例中,如果作业连续抛出 10 个异常,我们将等待 5 分钟,然后再次尝试该作业,但受到 30 分钟时间限制的限制。

当作业抛出异常但尚未达到异常阈值时,通常会立即重试该作业。但是,你可以在将中间件附加到作业时通过调用 backoff 方法来指定此类作业应延迟的分钟数:

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

在内部,该中间件使用Laravel的缓存系统来实现速率限制,并利用作业的类名作为缓存「键」。将中间件附加到作业时,你可以通过调用 by 方法来覆盖此键。如果你有多个作业与同一个第三方服务交互,并且你希望它们共享一个公共的节流「桶」,以确保它们遵守单个共享限制,这可能会很有用:

use Illuminate\Queue\Middleware\ThrottlesExceptions;
php
/**
 * 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;
php
/**
 * 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
    )];
}

如果你希望将受限制的异常报告给应用程序的异常处理程序,则可以通过在将中间件附加到作业时调用 report 方法来实现。或者,你可以为 report 方法提供一个闭包,并且仅当给定的闭包返回 true 时才会报告异常:

use Illuminate\Http\Client\HttpClientException;
use Illuminate\Queue\Middleware\ThrottlesExceptions;
php
/**
 * 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
    )];
}

INFO

如果你使用的是 Redis,则可以使用 Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis 中间件,该中间件针对 Redis 进行了微调,并且比基本异常节流中间件更高效:

跳过任务

Skip 中间件允许你指定应跳过/删除作业,而无需修改作业的逻辑。如果给定条件评估为 true,则 Skip::when 方法将删除作业,而如果条件评估为 false,则 Skip::unless 方法将删除作业:

use Illuminate\Queue\Middleware\Skip;
php
/**
* Get the middleware the job should pass through.
*/
public function middleware(): array
{
    return [
        Skip::when($someCondition),
    ];
}

你还可以将 Closure 传递给 whenunless 方法以进行更复杂的条件评估:

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

分发任务

编写作业类后,你可以使用作业本身的 dispatch 方法来分派它。传递给 dispatch 方法的参数将传递给作业的构造函数:

php
<?php

namespace App\Http\Controllers;

use App\Http\Controllers\Controller;
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');
    }
}

如果你想有条件地调度作业,你可以使用 dispatchIfdispatchUnless 方法:

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

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

在新的 Laravel 应用中,默认队列驱动是 sync。该驱动在当前请求的前台同步执行作业,本地开发时通常很方便。若希望真正将作业排队供后台处理,可在应用的 config/queue.php 配置文件中指定其他队列驱动。

延迟分发

如果你想指定作业不应立即可供队列工作人员处理,则可以在分派作业时使用 delay 方法。例如,我们指定作业在分派后 10 分钟后才能进行处理:

php
<?php

namespace App\Http\Controllers;

use App\Http\Controllers\Controller;
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()->addMinutes(10));

        return redirect('/podcasts');
    }
}

在某些情况下,作业可能配置了默认延迟。如果你需要绕过此延迟并分派作业立即处理,你可以使用 withoutDelay 方法:

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

WARNING

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

在响应发送到浏览器之后分派

或者,若 Web 服务器使用 FastCGI,dispatchAfterResponse 方法会延迟到 HTTP 响应发送到用户浏览器之后再分派作业。这样即使用户仍可开始使用应用,排队作业也仍在执行。通常仅应用于耗时约一秒的作业,例如发送邮件。由于这类作业在当前 HTTP 请求内处理,以此方式分派时无需运行队列 worker:

use App\Jobs\SendNotification;

SendNotification::dispatchAfterResponse();

你也可以 dispatch 一个闭包,并在 dispatch 辅助函数上链式调用 afterResponse 方法,以便在 HTTP 响应发送到浏览器之后执行该闭包:

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

dispatch(function () {
    Mail::to('taylor@example.com')->send(new WelcomeMessage);
})->afterResponse();

同步分发

如果你想立即(同步)调度作业,你可以使用 dispatchSync 方法。使用此方法时,作业不会排队,而是在当前进程内立即执行:

php
<?php

namespace App\Http\Controllers;

use App\Http\Controllers\Controller;
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');
    }
}

任务与数据库事务

虽然在数据库事务中分派作业是完全可以的,但你应该特别小心以确保你的作业实际上能够成功执行。在事务内分派作业时,该作业可能会在父事务提交之前由工作人员处理。发生这种情况时,你在数据库事务期间对模型或数据库记录所做的任何更新可能尚未反映在数据库中。此外,在事务中创建的任何模型或数据库记录可能不存在于数据库中。

值得庆幸的是,Laravel 提供了几种解决此问题的方法。首先,你可以在队列连接的配置数组中设置 after_commit 连接选项:

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

after_commit选项为true时,可以在数据库事务中调度作业;但是,Laravel 将等到开放父数据库事务提交后才实际分派作业。当然,如果当前没有数据库事务打开,则将立即分派作业。

如果事务由于事务期间发生异常而回滚,则该事务期间分派的作业将被丢弃。

INFO

after_commit 配置选项设置为 true 还会导致在提交所有打开的数据库事务后调度任何排队的事件侦听器、邮件、通知和广播事件。

内联指定提交后分发行为

如果你未将 after_commit 队列连接配置选项设置为 true,你仍然可以指示在提交所有打开的数据库事务后应分派特定作业。为了实现这一点,你可以将 afterCommit 方法链接到你的调度操作上:

use App\Jobs\ProcessPodcast;

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

同样,如果 after_commit 配置选项设置为 true,你可以指示应立即分派特定作业,而无需等待任何打开的数据库事务提交:

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

任务链

作业链允许你指定在主作业成功执行后应按顺序运行的排队任务列表。如果序列中的一项作业失败,则其余作业将不会运行。要执行排队任务链,你可以使用 Bus 门面提供的 chain 方法。 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();

WARNING

在作业中使用 $this->delete() 方法删除作业不会阻止处理链接作业。仅当链中的作业失败时,链才会停止执行。

链式连接与队列

如果你想指定用于链接作业的连接和队列,你可以使用 onConnectiononQueue 方法。这些方法指定应使用的队列连接和队列名称,除非已为排队任务显式分配不同的连接/队列:

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

向链中添加任务

有时,你可能需要将某个作业从现有作业链中的另一个作业添加到该链中。你可以使用 prependToChainappendToChain 方法来完成此操作:

php
/**
 * 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();

WARNING

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

自定义队列与连接

分发到指定队列

通过将作业推送到不同的队列,你可以对排队的作业进行「分类」,甚至可以优先考虑分配给各个队列的工作人员数量。请记住,这不会将作业推送到队列配置文件定义的不同队列「连接」,而只会推送到单个连接中的特定队列。要指定队列,请在分派作业时使用 onQueue 方法:

php
<?php

namespace App\Http\Controllers;

use App\Http\Controllers\Controller;
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
<?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
<?php

namespace App\Http\Controllers;

use App\Http\Controllers\Controller;
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');
    }
}

你可以将 onConnectiononQueue 方法链接在一起来指定作业的连接和队列:

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

或者,你可以通过在作业的构造函数中调用 onConnection 方法来指定作业的连接:

php
<?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');
    }
}

指定最大尝试次数 / 超时时间

最大尝试次数

若某个排队作业出错,你通常不希望它无限重试。因此,Laravel 提供了多种方式来指定作业可尝试的次数或时长。

指定作业可以尝试的最大次数的一种方法是通过 Artisan 命令行上的 --tries 开关。这将适用于工作人员处理的所有作业,除非正在处理的作业指定了可以尝试的次数:

shell
php artisan queue:work --tries=3

如果作业超过其最大尝试次数,它将被视为「失败」作业。有关处理失败作业的更多信息,请参阅 failed job documentation。如果向 queue:work 命令提供 --tries=0,则作业将无限期重试。

你也可以在任务类本身上定义最大尝试次数,以获得更细粒度的控制。若任务上指定了最大尝试次数,将优先于命令行提供的 --tries 值:

php
<?php

namespace App\Jobs;

class ProcessPodcast implements ShouldQueue
{
    /**
     * The number of times the job may be attempted.
     *
     * @var int
     */
    public $tries = 5;
}

如果你需要动态控制特定作业的最大尝试次数,你可以在作业上定义 tries 方法:

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

基于时间的尝试

作为定义作业失败之前可以尝试多少次的替代方法,你可以定义不应再尝试该作业的时间。这允许在给定时间范围内尝试任意次数的作业。要定义不应再尝试作业的时间,请将 retryUntil 方法添加到作业类中。此方法应返回 DateTime 实例:

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

INFO

你也可以在排队事件监听器上定义 tries 属性或 retryUntil 方法。

最大异常次数

有时你可能希望任务可尝试多次,但若重试由给定数量的未处理异常触发(而非直接由 release 方法释放),则应失败。为此,可在任务类上定义 maxExceptions 属性:

php
<?php

namespace App\Jobs;

use Illuminate\Support\Facades\Redis;

class ProcessPodcast implements ShouldQueue
{
    /**
     * The number of times the job may be attempted.
     *
     * @var int
     */
    public $tries = 25;

    /**
     * The maximum number of unhandled exceptions to allow before failing.
     *
     * @var int
     */
    public $maxExceptions = 3;

    /**
     * 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 锁,作业将被释放 10 秒,并且将继续重试最多 25 次。但是,如果作业抛出三个未处理的异常,作业将会失败。

超时

通常,你大致知道排队任务预计需要多长时间。因此,Laravel 允许你指定「超时」值。默认情况下,超时值为 60 秒。如果作业的处理时间超过超时值指定的秒数,则处理该作业的工作程序将退出并出现错误。通常,worker 将由 process manager configured on your server 自动重新启动。

作业可以运行的最大秒数可以使用 Artisan 命令行上的 --timeout 开关指定:

shell
php artisan queue:work --timeout=30

如果作业因持续超时而超出其最大尝试次数,则会将其标记为失败。

你也可以在任务类本身上定义允许任务运行的最大秒数。若任务上指定了超时,将优先于命令行上的任何超时设置:

php
<?php

namespace App\Jobs;

class ProcessPodcast implements ShouldQueue
{
    /**
     * The number of seconds the job can run before timing out.
     *
     * @var int
     */
    public $timeout = 120;
}

有时,IO 阻塞进程(例如套接字或传出 HTTP 连接)可能不遵守你指定的超时。因此,在使用这些功能时,你应该始终尝试使用它们的 API 来指定超时。例如,当使用 Guzzle 时,你应该始终指定连接和请求超时值。

WARNING

必须安装 PCNTL PHP 扩展才能指定任务超时。此外,任务的「timeout」值应始终小于其「retry after」 值;否则任务可能在实际执行完成或超时之前被再次尝试。

超时时标记失败

若希望任务在超时时标记为失败,可在任务类上定义 $failOnTimeout 属性:

php
/**
 * Indicate if the job should be marked as failed on timeout.
 *
 * @var bool
 */
public $failOnTimeout = true;

错误处理

如果在处理作业时引发异常,该作业将自动释放回队列,以便可以再次尝试。该作业将继续被释放,直到尝试次数达到你的应用程序允许的最大次数为止。最大尝试次数由 queue:work Artisan 命令上使用的 --tries 开关定义。或者,可以在作业类别本身上定义最大尝试次数。有关运行队列工作程序 can be found below 的更多信息。

手动释放任务

有时你可能希望手动将作业释放回队列中,以便稍后可以再次尝试。你可以通过调用 release 方法来完成此操作:

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

    $this->release();
}

默认情况下,release 方法会将作业释放回队列以立即处理。但是,你可以通过将整数或日期实例传递给 release 方法来指示队列在给定的秒数过去之前不让作业可供处理:

$this->release(10);

$this->release(now()->addSeconds(10));

手动使任务失败

有时你可能需要手动将作业标记为「失败」。为此,你可以调用 fail 方法:

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

    $this->fail();
}

如果你想因为捕获到异常而将作业标记为失败,则可以将该异常传递给 fail 方法。或者,为了方便起见,你可以传递字符串错误消息,该消息将为你转换为异常:

$this->fail($exception);

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

INFO

有关失败作业的更多信息,请查看 documentation on dealing with job failures

任务批处理

Laravel 的作业批处理功能让你可以轻松执行一批作业,并在整批作业执行完成后执行某些操作。开始之前,应创建数据库迁移以构建一张表,用于存放作业批次的元信息(例如完成百分比)。可使用 make:queue-batches-table Artisan 命令生成该迁移:

shell
php artisan make:queue-batches-table

php artisan migrate

定义可批处理任务

要定义可批处理作业,你应该照常 create a queueable job;但是,你应该将 Illuminate\Bus\Batchable 特征添加到作业类中。此特征提供对 batch 方法的访问,该方法可用于检索作业正在执行的当前批处理:

php
<?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...
    }
}

分发批次

要分派一批作业,应使用 Bus 门面的 batch 方法。当然,批处理主要在与完成回调结合时才更有用。因此,你可以使用 thencatchfinally 方法为批次定义完成回调。每个回调被调用时都会收到一个 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) {
    // First batch job failure detected...
})->finally(function (Batch $batch) {
    // The batch has finished executing...
})->dispatch();

return $batch->id;

批次的 ID(可通过 $batch->id 属性访问)可用于 query the Laravel command bus 以获取有关批次分派后的信息。

WARNING

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

命名批次

如果批处理已命名,某些工具(例如 Laravel HorizonLaravel Telescope)可能会为批处理提供更用户友好的调试信息。要为批次分配任意名称,你可以在定义批次时调用 name 方法:

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

批次连接与队列

如果你想指定用于批处理作业的连接和队列,你可以使用 onConnectiononQueue 方法。所有批处理作业必须在同一连接和队列中执行:

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

链式与批次

你可以通过将链接的作业放置在数组中来在批处理中定义一组 chained jobs。例如,我们可以并行执行两个作业链,并在两个作业链都完成处理时执行回调:

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) {
    // ...
})->dispatch();

相反,你可以通过在链中定义批次来在 chain 内运行批量作业。例如,你可以首先运行一批作业来发布多个播客,然后运行一批作业来发送发布通知:

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;
php
/**
 * Execute the job.
 */
public function handle(): void
{
    if ($this->batch()->cancelled()) {
        return;
    }

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

WARNING

你只能将属于同一批次的作业中的作业添加到批次中。

检查批次

提供给批次完成回调的 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 有效负载。这样可以方便地在应用程序的 UI 中显示有关批次完成进度的信息。

要通过 ID 检索批次,你可以使用 Bus 门面的 findBatch 方法:

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

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

取消批次

有时你可能需要取消给定批次的执行。这可以通过在 Illuminate\Bus\Batch 实例上调用 cancel 方法来完成:

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

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

正如你在前面的示例中可能已经注意到的那样,批处理作业通常应在继续执行之前确定其相应的批处理是否已被取消。但是,为了方便起见,你可以将 SkipIfBatchCancelled middleware 分配给作业。顾名思义,如果相应的批次已被取消,该中间件将指示 Laravel 不处理该作业:

use Illuminate\Queue\Middleware\SkipIfBatchCancelled;
php
/**
 * 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();

重试失败的批次任务

为方便起见,Laravel 提供了 queue:retry-batch Artisan 命令,可轻松重试给定批次中的全部失败作业。该命令接受应重试失败作业的批次 UUID:

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

修剪批次

如果不进行修剪,job_batches 表可以非常快速地积累记录。为了缓解这种情况,你应该每天运行 schedule 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 表可能会累积从未成功完成的批次的批次记录,例如作业失败且该作业从未成功重试的批次。你可以使用 unfinished 选项指示 queue:prune-batches 命令修剪这些未完成的批次记录:

use Illuminate\Support\Facades\Schedule;

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

同样,你的 job_batches 表也可能会累积已取消批次的批次记录。你可以使用 cancelled 选项指示 queue:prune-batches 命令删除这些已取消的批次记录:

use Illuminate\Support\Facades\Schedule;

Schedule::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 应用程序的作业批次。

此外,如果你想利用 automatic batch pruning,你可以为表定义 ttl 属性。

DynamoDB 配置

接下来,安装 AWS 开发工具包,以便你的 Laravel 应用程序可以与 Amazon DynamoDB 通信:

shell
composer require aws/aws-sdk-php

然后,将 queue.batching.driver 配置选项的值设置为 dynamodb。此外,你还应该在 batching 配置数组中定义 keysecretregion 配置选项。这些选项将用于向 AWS 进行身份验证。使用 dynamodb 驱动程序时,不需要 queue.batching.database 配置选项:

php
'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's native TTL functionality 自动删除旧批次的记录。

如果你使用 ttl 属性定义 DynamoDB 表,你可以定义配置参数来指示 Laravel 如何修剪批次记录。 queue.batching.ttl_attribute 配置值定义保存 TTL 的属性的名称,而 queue.batching.ttl 配置值定义相对于上次更新记录的时间,可以从 DynamoDB 表中删除批记录的秒数:

php
'batching' => [
    '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' => 'job_batches',
    'ttl_attribute' => 'ttl',
    'ttl' => 60 * 60 * 24 * 7, // 7 days...
],

将闭包加入队列

你也可以分派一个闭包,而不是将作业类分派到队列。这对于需要在当前请求周期之外执行的快速、简单的任务非常有用。将闭包分派到队列时,闭包的代码内容经过加密签名,以便在传输过程中无法修改:

$podcast = App\Podcast::find(1);

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

使用 catch 方法,你可以提供一个闭包,如果排队的闭包在耗尽队列的所有 configured retry attempts 后未能成功完成,则应执行该闭包:

use Throwable;

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

WARNING

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

运行队列 Worker

queue:work 命令

Laravel 包括 Artisan 命令,该命令将启动队列工作程序并在新作业推送到队列时处理它们。你可以使用 queue:work Artisan 命令运行工作程序。请注意,一旦 queue:work 命令启动,它将继续运行,直到手动停止或你关闭终端:

shell
php artisan queue:work

INFO

要保持 queue:work 进程在后台永久运行,你应该使用进程监视器(例如 Supervisor)来确保队列工作程序不会停止运行。

若希望在命令输出中包含已处理的作业 ID,可在调用 queue:work 命令时加上 -v 标志:

shell
php artisan queue:work -v

请记住,队列工作进程是长期存在的进程,并将启动的应用程序状态存储在内存中。因此,他们在启动后不会注意到你的代码库中的更改。因此,在部署过程中,请务必restart your queue workers。此外,请记住,应用程序创建或修改的任何静态状态都不会在作业之间自动重置。

或者,你可以运行 queue:listen 命令。使用 queue:listen 命令时,当你想要重新加载更新的代码或重置应用程序状态时,无需手动重新启动工作程序;但是,此命令的效率明显低于 queue:work 命令:

shell
php artisan queue:listen

运行多个队列 Worker

要将多个工作人员分配到队列并同时处理作业,你只需启动多个 queue:work 进程即可。这可以通过终端中的多个选项卡在本地完成,也可以使用流程管理器的配置设置在生产中完成。 When using Supervisor,你可以使用 numprocs 配置值。

指定连接与队列

你还可以指定工作人员应使用哪个队列连接。传递给 work 命令的连接名称应对应于 config/queue.php 配置文件中定义的连接之一:

shell
php artisan queue:work redis

默认情况下,queue:work 命令仅处理给定连接上默认队列的作业。但是,你可以通过仅处理给定连接的特定队列来进一步自定义队列工作程序。例如,如果你的所有电子邮件都在 redis 队列连接上的 emails 队列中处理,你可以发出以下命令来启动仅处理该队列的工作进程:

shell
php artisan queue:work redis --queue=emails

处理指定数量的任务

--once 选项可用于指示工作线程仅处理队列中的单个作业:

shell
php artisan queue:work --once

--max-jobs 选项可用于指示工作线程处理给定数量的作业,然后退出。当与 Supervisor 结合使用时,此选项可能很有用,以便你的工作线程在处理给定数量的作业后自动重新启动,释放它们可能积累的任何内存:

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

处理完所有排队任务后退出

--stop-when-empty 选项可用于指示工作线程处理所有作业,然后正常退出。如果你希望在队列为空后关闭容器,则在处理 Docker 容器内的 Laravel 队列时,此选项非常有用:

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

在给定秒数内处理任务

--max-time 选项可用于指示工作线程在给定的秒数内处理作业,然后退出。当与 Supervisor 结合使用时,此选项可能很有用,以便你的工作线程在处理给定时间的作业后自动重新启动,释放它们可能积累的任何内存:

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

Worker 休眠时长

当队列中有可用作业时,工作人员将继续处理作业,作业之间不会有任何延迟。但是,sleep 选项确定在没有可用作业的情况下工作人员将「睡眠」的秒数。当然,在睡眠时,worker不会处理任何新的作业:

shell
php artisan queue:work --sleep=3

维护模式与队列

当你的应用程序位于 maintenance mode 时,不会处理任何排队的作业。一旦应用程序退出维护模式,作业将继续正常处理。

要强制队列工作人员在启用维护模式的情况下处理作业,你可以使用 --force 选项:

shell
php artisan queue:work --force

资源考量

守护进程队列 worker 在处理每个任务前不会「重启」框架。因此,应在每个任务完成后释放沉重资源。例如,若使用 GD 库 处理图像,处理完成后应使用 imagedestroy 释放内存。

队列优先级

有时你可能希望优先考虑队列的处理方式。例如,在 config/queue.php 配置文件中,你可以将 redis 连接的默认 queue 设置为 low。但是,有时你可能希望将作业推送到 high 优先级队列,如下所示:

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

要启动一个工作程序,在继续处理 low 队列上的任何作业之前验证是否已处理所有 high 队列作业,请将逗号分隔的队列名称列表传递给 work 命令:

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

队列 Worker 与部署

由于队列工作进程是长期存在的进程,因此如果不重新启动,它们将不会注意到代码的更改。因此,使用队列 Worker部署应用程序的最简单方法是在部署过程中重新启动工作线程。你可以通过发出 queue:restart 命令正常重启所有工作进程:

shell
php artisan queue:restart

此命令将指示所有队列工作人员在处理完当前作业后正常退出,以便不会丢失现有作业。由于执行 queue:restart 命令时队列工作程序将退出,因此你应该运行进程管理器(例如 Supervisor)来自动重新启动队列工作程序。

INFO

队列使用 cache 来存储重新启动信号,因此在使用此功能之前,你应该验证是否为你的应用程序正确配置了缓存驱动程序。

任务过期与超时

任务过期

config/queue.php 配置文件中,每个队列连接定义一个 retry_after 选项。此选项指定队列连接在重试正在处理的作业之前应等待的秒数。例如,如果将 retry_after 的值设置为 90,则如果作业已处理 90 秒而没有被释放或删除,则该作业将被释放回队列中。通常,你应该将 retry_after 值设置为作业完成处理所需合理时间的最大秒数。

WARNING

唯一不包含 retry_after 值的队列连接是 Amazon SQS。 SQS 将根据 AWS 控制台中管理的 Default Visibility Timeout 重试作业。

Worker 超时

queue:work Artisan 命令公开 --timeout 选项。默认情况下,--timeout 值为 60 秒。如果作业的处理时间超过超时值指定的秒数,则处理该作业的工作程序将退出并出现错误。通常,worker 将由 process manager configured on your server 自动重新启动:

shell
php artisan queue:work --timeout=60

retry_after 配置选项和 --timeout CLI 选项不同,但协同工作以确保作业不会丢失并且作业仅成功处理一次。

WARNING

--timeout 值应始终比 retry_after 配置值至少短几秒。这将确保处理冻结作业的工作人员始终在重试作业之前终止。如果你的 --timeout 选项比 retry_after 配置值长,你的作业可能会被处理两次。

Supervisor 配置

在生产中,你需要一种方法来保持 queue:work 进程运行。 queue:work 进程可能会因多种原因停止运行,例如超出工作超时或执行 queue:restart 命令。

因此,你需要配置一个进程监视器,它可以检测 queue:work 进程何时退出并自动重新启动它们。此外,进程监视器可以允许你指定要同时运行的 queue:work 进程数。 Supervisor是Linux环境中常用的进程监视器,我们将在下面的文档中讨论如何配置它。

安装 Supervisor

Supervisor 是 Linux 操作系统的进程监视器,如果 queue:work 进程失败,它将自动重新启动它们。要在 Ubuntu 上安装 Supervisor,你可以使用以下命令:

shell
sudo apt-get install supervisor

INFO

若自行配置与管理 Supervisor 让你感到吃力,可考虑使用 Laravel Forge,它会为你的生产 Laravel 项目自动安装并配置 Supervisor。

配置 Supervisor

Supervisor 配置文件通常存储在 /etc/supervisor/conf.d 目录中。在此目录中,你可以创建任意数量的配置文件来指示主管如何监视你的进程。例如,让我们创建一个启动并监视 queue:work 进程的 laravel-worker.conf 文件:

ini
[program:laravel-worker]
process_name=%(program_name)s_%(process_num)02d
command=php /home/forge/app.com/artisan queue:work sqs --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 指令以反映你所需的队列连接和工作选项。

WARNING

你应该确保 stopwaitsecs 的值大于运行时间最长的作业所消耗的秒数。否则,Supervisor 可能会在处理完成之前终止作业。

启动 Supervisor

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

shell
sudo supervisorctl reread

sudo supervisorctl update

sudo supervisorctl start "laravel-worker:*"

有关 Supervisor 的更多信息,请参阅 Supervisor documentation

处理失败任务

有时你的排队任务会失败。别担心,事情并不总是按计划进行! Laravel 包括一种方便的方法来 specify the maximum number of times a job should be attempted。异步作业超过此尝试次数后,将被插入到 failed_jobs 数据库表中。失败的 Synchronously dispatched jobs 不会存储在此表中,并且它们的异常会立即由应用程序处理。

用于创建 failed_jobs 表的迁移通常已存在于新的 Laravel 应用程序中。但是,如果你的应用程序不包含此表的迁移,你可以使用 make:queue-failed-table 命令来创建迁移:

shell
php artisan make:queue-failed-table

php artisan migrate

运行队列 worker 时,可通过 queue:work 命令的 --tries 开关指定任务最大尝试次数。若未指定 --tries,任务将只尝试一次,或按任务类 $tries 属性指定的次数尝试:

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

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

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

若希望按任务配置 Laravel 在遇到异常后重试前等待的秒数,可在任务类上定义 backoff 属性:

php
/**
 * The number of seconds to wait before retrying the job.
 *
 * @var int
 */
public $backoff = 3;

如果你需要更复杂的逻辑来确定作业的退避时间,你可以在作业类上定义 backoff 方法:

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

可通过让 backoff 方法返回退避值数组,轻松配置「指数」退避。本例中,第一次重试延迟 1 秒,第二次 5 秒,第三次 10 秒;若仍有剩余尝试次数,之后每次重试均为 10 秒:

php
/**
* Calculate the number of seconds to wait before retrying the job.
*
* @return array<int, int>
*/
public function backoff(): array
{
    return [1, 5, 10];
}

失败任务后的清理

当特定作业失败时,你可能希望向用户发送警报或恢复该作业部分完成的任何操作。为此,你可以在作业类上定义 failed 方法。导致作业失败的 Throwable 实例将传递给 failed 方法:

php
<?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...
    }
}

WARNING

在调用 failed 方法之前实例化作业的新实例;因此,在 handle 方法中可能发生的任何类属性修改都将丢失。

重试失败任务

要查看已插入 failed_jobs 数据库表中的所有失败作业,你可以使用 queue:failed Artisan 命令:

shell
php artisan queue:failed

queue:failed 命令将列出作业 ID、连接、队列、失败时间以及有关作业的其他信息。作业 ID 可用于重试失败的作业。例如,要重试 ID 为 ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece 的失败作业,请发出以下命令:

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

如有必要,你可以将多个 ID 传递给命令:

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

你还可以重试特定队列的所有失败作业:

shell
php artisan queue:retry --queue=name

要重试所有失败的作业,请执行 queue:retry 命令并传递 all 作为 ID:

shell
php artisan queue:retry all

如果你想删除失败的作业,你可以使用 queue:forget 命令:

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

INFO

使用 Horizon 时,应使用 horizon:forget 命令来删除失败的作业,而不是使用 queue:forget 命令。

要从 failed_jobs 表中删除所有失败的作业,你可以使用 queue:flush 命令:

shell
php artisan queue:flush

忽略缺失模型

将 Eloquent 模型注入作业时,该模型会自动序列化,然后放入队列中,并在处理作业时从数据库中重新检索。但是,如果在作业等待工作人员处理时模型已被删除,你的作业可能会失败并显示 ModelNotFoundException

为方便起见,可将任务的 deleteWhenMissingModels 属性设为 true,以自动删除缺少模型的任务。设为 true 时,Laravel 会静默丢弃该任务且不抛出异常:

php
/**
 * Delete the job if its models no longer exist.
 *
 * @var bool
 */
public $deleteWhenMissingModels = true;

修剪失败任务

你可以通过调用 queue:prune-failed Artisan 命令来修剪应用程序的 failed_jobs 表中的记录:

shell
php artisan queue:prune-failed

默认情况下,所有超过 24 小时的失败作业记录都将被删除。如果为该命令提供 --hours 选项,则仅保留最近 N 小时内插入的失败作业记录。例如,以下命令将删除超过 48 小时前插入的所有失败作业记录:

shell
php 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 开发工具包,以便你的 Laravel 应用程序可以与 Amazon DynamoDB 通信:

shell
composer require aws/aws-sdk-php

接下来,将 queue.failed.driver 配置选项的值设置为 dynamodb。此外,你应该在失败的作业配置数组中定义 keysecretregion 配置选项。这些选项将用于向 AWS 进行身份验证。使用 dynamodb 驱动程序时,不需要 queue.failed.database 配置选项:

php
'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 环境变量来完成:

ini
QUEUE_FAILED_DRIVER=null

失败任务事件

如果你想注册一个在作业失败时调用的事件侦听器,你可以使用 Queue 门面的 failing 方法。例如,我们可以从 Laravel 中包含的 AppServiceProviderboot 方法为此事件附加一个闭包:

php
<?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
        });
    }
}

从队列清除任务

INFO

使用 Horizon 时,应使用 horizon:clear 命令从队列中清除作业,而不是使用 queue:clear 命令。

如果你想从默认连接的默认队列中删除所有作业,你可以使用 queue:clear Artisan 命令来执行此操作:

shell
php artisan queue:clear

你还可以提供 connection 参数和 queue 选项来从特定连接和队列中删除作业:

shell
php artisan queue:clear redis --queue=emails

WARNING

从队列中清除作业仅适用于 SQS、Redis 和数据库队列驱动程序。此外,SQS 消息删除过程最多需要 60 秒,因此在清除队列后 60 秒内发送到 SQS 队列的作业也可能会被删除。

监控队列

如果你的队列突然收到大量作业,它可能会变得不堪重负,从而导致作业完成需要很长时间的等待。如果你愿意,Laravel 可以在队列作业计数超过指定阈值时向你发出警报。

首先,你应该将 queue:monitor 命令安排到 run every minute。该命令接受你希望监视的队列的名称以及所需的作业计数阈值:

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

单独安排此命令不足以触发通知,提醒你队列的不堪重负状态。当命令遇到作业计数超过阈值的队列时,将调度 Illuminate\Queue\Events\QueueBusy 事件。你可以在应用程序的 AppServiceProvider 中监听此事件,以便向你或你的开发团队发送通知:

php
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', 'dev@example.com')
            ->notify(new QueueHasLongWaitTime(
                $event->connection,
                $event->queue,
                $event->size
            ));
    });
}

测试

在测试分派作业的代码时,你可能希望指示 Laravel 实际上不执行作业本身,因为可以直接测试作业的代码,并与分派它的代码分开进行测试。当然,为了测试作业本身,你可以实例化一个作业实例并直接在测试中调用 handle 方法。

你可以使用 Queue 门面的 fake 方法来防止排队的作业实际被推送到队列中。调用 Queue 门面的 fake 方法后,你可以断言应用程序尝试将作业推送到队列:

php
<?php

use App\Jobs\AnotherJob;
use App\Jobs\FinalJob;
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 twice...
    Queue::assertPushed(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 the total number of jobs that were pushed...
    Queue::assertCount(3);
});
php
<?php

namespace Tests\Feature;

use App\Jobs\AnotherJob;
use App\Jobs\FinalJob;
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 twice...
        Queue::assertPushed(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 the total number of jobs that were pushed...
        Queue::assertCount(3);
    }
}

你可以将闭包传递给 assertPushedassertNotPushed 方法,以断言已推送的作业通过了给定的「真值测试」。若至少有一个作业通过该真值测试,则断言成功:

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

伪造部分任务

如果你只需要伪造特定作业,同时允许其他作业正常执行,你可以将应伪造的作业的类名传递给 fake 方法:

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

    // Perform order shipping...

    // Assert a job was pushed twice...
    Queue::assertPushed(ShipOrder::class, 2);
});
php
public function test_orders_can_be_shipped(): void
{
    Queue::fake([
        ShipOrder::class,
    ]);

    // Perform order shipping...

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

你可以使用 except 方法伪造除一组指定作业之外的所有作业:

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

测试任务链

要测试作业链,你需要利用 Bus 外观的伪造功能。 Bus 门面的 assertChained 方法可用于断言 chain of jobs 已调度。 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);

测试链式修改

如果链式作业 prepends or appends jobs to an existing chain,你可以使用该作业的 assertHasChain 方法来断言该作业具有预期的剩余作业链:

php
$job = new ProcessPodcast;

$job->handle();

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

assertDoesntHaveChain 方法可用于断言作业的剩余链为空:

php
$job->assertDoesntHaveChain();

测试链式批次

如果你的作业链 contains a batch of jobs,你可以通过在链断言中插入 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 方法可用于断言 batch of jobs 已调度。给予 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;
});

你可以使用 assertBatchCount 方法来断言已调度给定数量的批次:

Bus::assertBatchCount(3);

你可以使用 assertNothingBatched 来断言没有发送任何批次:

Bus::assertNothingBatched();

测试任务 / 批次交互

此外,你有时可能需要测试单个作业与其基础批次的交互。例如,你可能需要测试作业是否取消了其批次的进一步处理。为此,你需要通过 withFakeBatch 方法为作业分配一个假批次。 withFakeBatch 方法返回一个包含作业实例和假批次的元组:

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

$job->handle();

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

测试任务 / 队列交互

有时,你可能需要测试排队任务 releases itself back onto the queue。或者,你可能需要测试作业是否已自行删除。你可以通过实例化作业并调用 withFakeQueueInteractions 方法来测试这些队列交互。

在作业的队列交互被伪造后,可调用作业上的 handle 方法。调用后,可使用 assertReleasedassertDeletedassertNotDeletedassertFailedassertFailedWithassertNotFailed 方法对作业的队列交互进行断言:

php
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();

任务事件

Queue facade 上使用 beforeafter 方法,你可以指定在处理排队任务之前或之后执行的回调。这些回调是为仪表板执行额外日志记录或增量统计的绝佳机会。通常,你应该从 服务提供者boot 方法调用这些方法。例如,我们可以使用 Laravel 附带的 AppServiceProvider

php
<?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()
        });
    }
}

Queue facade 上使用 looping 方法,你可以指定在工作线程尝试从队列中获取作业之前执行的回调。例如,你可以注册一个闭包来回滚之前失败的作业留下的任何事务:

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

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