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 属性中定义的队列上:

php
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 键都放入同一个哈希槽中:

php
'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 来指示驱动程序在等待作业可用时应阻塞五秒钟:

php
'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: 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 方法依赖注入

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

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

php
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 属性:

php
use Illuminate\Queue\Attributes\WithoutRelations;

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

为了方便起见,如果你希望序列化所有没有关系的模型,你可以将 WithoutRelations 属性应用于整个类,而不是将该属性应用于每个模型:

php
<?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 模型的集合或数组而不是单个模型,则在反序列化和执行作业时,该集合中的模型将不会恢复其关系。这是为了防止处理大量模型的作业过度使用资源。

唯一任务

WARNING

独特的作业需要支持 locks 的缓存驱动程序。目前,memcachedredisdynamodbdatabasefilearray 缓存驱动程序支持原子锁。

WARNING

独特的作业限制不适用于批次内的作业。

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

php
<?php

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

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

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

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

php
<?php

namespace App\Jobs;

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

class UpdateSearchIndex implements ShouldQueue, ShouldBeUnique
{
    /**
     * The product instance.
     *
     * @var \App\Models\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 Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Contracts\Queue\ShouldBeUniqueUntilProcessing;

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

唯一任务锁

在后台,当调度 ShouldBeUnique 作业时,Laravel 尝试使用 uniqueId 键获取 lock。如果已持有锁,则不会分派作业。当作业完成处理或所有重试尝试失败时,将释放此锁。默认情况下,Laravel将使用默认的缓存驱动程序来获取此锁。但是,如果你希望使用其他驱动程序来获取锁,则可以定义一个 uniqueVia 方法,该方法返回应使用的缓存驱动程序:

php
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 速率限制功能,仅允许每五秒处理一个作业:

php
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
<?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 一样,作业中间件接收正在处理的作业以及应调用以继续处理作业的回调。

你可以使用 make:job-middleware Artisan 命令生成新的作业中间件类。创建作业中间件后,可以通过从作业的 middleware 方法返回它们来将它们附加到作业。由 make:job Artisan 命令搭建的作业中不存在此方法,因此你需要手动将其添加到作业类中:

php
use App\Jobs\Middleware\RateLimited;

/**
 * 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

php
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方法;但是,该值最常用于按客户划分速率限制:

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

定义速率限制后,你可以使用 Illuminate\Queue\Middleware\RateLimited 中间件将速率限制器附加到你的作业。每次作业超过速率限制时,该中间件都会根据速率限制持续时间以适当的延迟将作业释放回队列:

php
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 总数。你可能需要相应调整任务类上的 triesmaxExceptions 属性;也可以使用 retryUntil 方法 定义任务不再被尝试的截止时间。

使用 releaseAfter 方法,你还可以指定再次尝试释放的作业之前必须经过的秒数:

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

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

php
/**
 * 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 进行了微调,并且比基本限速中间件更高效:

php
use Illuminate\Queue\Middleware\RateLimitedWithRedis;

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

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

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

防止任务重叠

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

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

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

将重叠任务释放回队列仍会增加该任务的总尝试次数。你可能需要相应调整任务类上的 triesmaxExceptions。例如,若保持默认的 tries 为 1,重叠任务稍后将无法重试。

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

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 的作业配对:

php
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 方法来指定此类作业应延迟的分钟数:

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

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

php
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 时,异常才会被限制:

php
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 方法允许你在发生给定异常时完全删除作业:

php
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 时才会报告异常:

php
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 进行了微调,并且比基本异常节流中间件更高效:

php
use Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis;

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

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

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

跳过任务

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

php
use Illuminate\Queue\Middleware\Skip;

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

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

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

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

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

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

在新的 Laravel 应用程序中,database 连接被定义为默认队列。你可以通过更改应用程序的 .env 文件中的 QUEUE_CONNECTION 环境变量来指定不同的默认队列连接。

延迟分发

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

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

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

WARNING

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

同步分发

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

php
<?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 连接:

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

deferred 连接也用作默认 failover queue

同样,background 连接在 HTTP 响应发送给用户后处理作业;但是,该作业是在单独生成的 PHP 进程中处理的,从而允许 PHP-FPM / 应用程序工作线程可用于处理另一个传入的 HTTP 请求:

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

任务与数据库事务

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

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

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

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

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

INFO

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

内联指定提交后分发行为

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

php
use App\Jobs\ProcessPodcast;

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

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

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

任务链

作业链允许你指定在主作业成功执行后应按顺序运行的排队任务列表。如果序列中的一项作业失败,则其余作业将不会运行。要执行排队任务链,你可以使用 Bus 门面提供的 chain 方法。 Laravel 的命令总线是一个较低级别的组件,队列作业调度构建在以下组件之上:

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

除了链接作业类实例之外,你还可以链接闭包:

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

WARNING

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

链式连接与队列

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

php
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 实例:

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

WARNING

通过构造函数中的 onQueue 指定队列仅对任务类有效。对于队列事件监听器,请在监听器类上定义 viaQueue 方法或 $queue 属性。

分发到指定连接

如果你的应用程序与多个队列连接交互,你可以使用 onConnection 方法指定将作业推送到哪个连接:

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

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

php
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 队列系统的核心概念,并为许多高级功能提供支持。虽然它们一开始看起来可能令人困惑,但在修改默认配置之前了解它们的工作原理非常重要。

当作业被调度时,它被推送到队列中。然后,一名工人捡起它并尝试执行它。这是一次工作尝试。

但是,尝试并不一定意味着作业的 handle 方法已执行。尝试也可以通过多种方式「消耗」:

  • 作业在执行期间遇到未处理的异常。
  • 使用 `$this->release()` 手动将作业释放回队列。
  • `WithoutOverlapping`、`RateLimited`等中间件获取锁失败,释放作业。
  • 作业超时。
  • 作业的 `handle` 方法运行并完成,不会引发异常。

你可能不想无限期地继续尝试某项工作。因此,Laravel 提供了各种方法来指定可以尝试作业的次数或时间。

INFO

默认情况下,Laravel 只会尝试一次作业。如果你的作业使用 WithoutOverlappingRateLimited 等中间件,或者你手动释放作业,则可能需要通过 tries 选项增加允许的尝试次数。

指定作业可以尝试的最大次数的一种方法是通过 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 实例:

php
use DateTime;

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

如果同时定义了 retryUntiltries,则 Laravel 优先于 retryUntil 方法。

INFO

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

最大异常次数

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

php
<?php

namespace App\Jobs;

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

class ProcessPodcast implements ShouldQueue
{
    use Queueable;

    /**
     * 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;

INFO

默认情况下,当作业超时时,它会消耗一次尝试并被释放回队列(如果允许重试)。但是,如果你将作业配置为超时失败,则无论尝试设置的值如何,都不会重试。

SQS FIFO 与公平队列

Laravel 支持 Amazon SQS FIFO(先进先出) 队列,可按发送的确切顺序处理任务,并通过消息去重确保恰好一次处理。

FIFO 队列需要消息组 ID 来确定哪些作业可以并行处理。具有相同组ID的作业是顺序处理的,而具有不同组ID的消息可以并发处理。

Laravel 提供了一个流畅的 onGroup 方法来在调度作业时指定消息组 ID:

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

SQS FIFO 队列支持消息重复数据删除,以确保一次性处理。在作业类中实现 deduplicationId 方法以提供自定义重复数据删除 ID:

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

FIFO 监听器、邮件与通知

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

要定义 queued event listener 的消息组,请在侦听器上定义 messageGroup 方法。你还可以选择定义 deduplicationId 方法:

php
<?php

namespace App\Listeners;

class SendShipmentNotification
{
    // ...

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

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

当发送将在 FIFO 队列中排队的 mail message 时,你应该在发送通知时调用 onGroup 方法,也可以选择调用 withDeduplicator 方法:

php
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 队列中排队的 notification 时,你应该在发送通知时调用 onGroup 方法,也可以选择调用 withDeduplicator 方法:

php
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 配置文件中包含示例故障转移配置:

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

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

ini
QUEUE_CONNECTION=failover

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

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

INFO

你不需要使用 syncbackgrounddeferred 队列驱动程序运行连接工作线程,因为这些驱动程序在当前 PHP 进程中处理作业。

当队列连接操作失败并激活故障转移时,Laravel 将调度 Illuminate\Queue\Events\QueueFailedOver 事件,允许你报告或记录队列连接失败。

INFO

如果你使用 Laravel Horizon,请记住 Horizon 仅管理 Redis 队列。如果你的故障转移列表包括 database,则你应该与 Horizon 一起运行常规 php artisan queue:work database 进程。

错误处理

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

手动释放任务

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

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

    $this->release();
}

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

php
$this->release(10);

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

手动使任务失败

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

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

    $this->fail();
}

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

php
$this->fail($exception);

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

INFO

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

在特定异常时使任务失败

FailOnException job middleware 允许你在引发特定异常时短路重试。这允许重试暂时性异常(例如外部 API 错误),但在持久性异常(例如用户权限被撤销)上永久失败作业:

php
<?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\Middleware\FailOnException;
use Illuminate\Support\Facades\Http;

class SyncChatHistory implements ShouldQueue
{
    use Queueable;

    public $tries = 3;

    /**
     * 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 命令生成此迁移:

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 实例。

当运行多个队列工作程序时,批处理中的作业将并行处理。因此,作业完成的顺序可能与它们添加到批次中的顺序不同。有关如何按顺序运行一系列作业的信息,请参阅有关 job chains and batches 的文档。

在此示例中,我们假设正在对一批作业进行排队,每个作业处理 CSV 文件中给定数量的行:

php
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 属性访问)可用于 query the Laravel command bus 以获取有关批次分派后的信息。

WARNING

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

命名批次

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

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

批次连接与队列

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

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

链式与批次

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

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

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

php
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 请求期间可能需要很长时间才能分派)时,此模式非常有用。因此,你可能希望分派第一批「加载器」作业,以便用更多作业来补充该批作业:

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

在此示例中,我们将使用 LoadImportBatch 作业将批次与其他作业混合。为了实现这一点,我们可以在批处理实例上使用 add 方法,该实例可以通过作业的 batch 方法访问:

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

WARNING

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

检查批次

提供给批次完成回调的 Illuminate\Bus\Batch 实例具有多种属性和方法,可帮助你与给定批次的作业进行交互和检查:

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

php
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()) {
        $this->batch()->cancel();

        return;
    }

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

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

php
use Illuminate\Queue\Middleware\SkipIfBatchCancelled;

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

批次失败

当批处理作业失败时,将调用 catch 回调(如果已分配)。仅针对批处理中失败的第一个作业调用此回调。

允许失败

当批次内的作业失败时,Laravel 会自动将该批次标记为「已取消」。如果你愿意,你可以禁用此行为,以便作业失败时不会自动将批次标记为已取消。这可以通过在分派批处理时调用 allowFailures 方法来完成:

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

你可以选择提供 allowFailures 方法的闭包,该方法将在每次作业失败时执行:

php
$batch = Bus::batch([
    // ...
])->allowFailures(function (Batch $batch, $exception) {
    // Handle individual job failures...
})->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 命令:

php
use Illuminate\Support\Facades\Schedule;

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

默认情况下,所有超过 24 小时的已完成批次都将被修剪。调用命令时可以使用 hours 选项来确定批量数据保留多长时间。例如,以下命令将删除 48 小时前完成的所有批次:

php
use Illuminate\Support\Facades\Schedule;

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

有时,你的 job_batches 表可能会累积从未成功完成的批次的批次记录,例如作业失败且该作业从未成功重试的批次。你可以使用 unfinished 选项指示 queue:prune-batches 命令修剪这些未完成的批次记录:

php
use Illuminate\Support\Facades\Schedule;

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

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

php
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...
],

将闭包加入队列

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

php
use App\Models\Podcast;

$podcast = Podcast::find(1);

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

要为排队的闭包分配一个名称,该名称可以由队列报告仪表板使用,并通过 queue:work 命令显示,你可以使用 name 方法:

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

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

php
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... {#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 优先级队列,如下所示:

php
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 配置值长,你的作业可能会被处理两次。

暂停与恢复队列 Worker

有时,你可能需要暂时阻止队列 Worker处理新作业,而不完全停止该工作线程。例如,你可能希望在系统维护期间暂停作业处理。 Laravel 提供 queue:pausequeue:continue Artisan 命令来暂停和恢复队列工作程序。

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

shell
php artisan queue:pause database:default

在此示例中,database 是队列连接名称,default 是队列名称。一旦队列暂停,处理该队列中作业的任何工作人员将继续完成当前作业,但在队列恢复之前不会拾取任何新作业。

要恢复处理暂停队列上的作业,请使用 queue:continue 命令:

shell
php artisan queue:continue database:default

恢复队列后,worker 将立即开始处理该队列中的新任务。请注意,暂停队列不会停止 worker 进程本身——仅阻止其处理指定队列中的新任务。

Worker 重启与暂停信号

默认情况下,队列工作人员会在每次作业迭代时轮询缓存驱动程序以获取重新启动和暂停信号。虽然此轮询对于响应 queue:restartqueue:pause 命令至关重要,但它确实会带来一点性能开销。

如果你需要优化性能并且不需要这些中断功能,你可以通过调用 Queue 外观上的 withoutInterruptionPolling 方法来全局禁用此轮询。这通常应该在 AppServiceProviderboot 方法中完成:

php
use Illuminate\Support\Facades\Queue;

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

或者,你可以通过在 Illuminate\Queue\Worker 类上设置静态 $restartable$pausable 属性来单独禁用重新启动或暂停轮询:

php
use Illuminate\Queue\Worker;

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

WARNING

当中断轮询被禁用时,工作人员将不会响应 queue:restartqueue:pause 命令(取决于禁用的功能)。

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 Cloud,它为运行 Laravel 队列工作程序提供了完全托管的平台。

配置 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 --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 方法。但是,如果作业因已达到允许的最大尝试次数而失败,则 $exception 将是 Illuminate\Queue\MaxAttemptsExceededException 的实例。同样,如果作业由于超过配置的超时而失败,则 $exception 将成为 Illuminate\Queue\TimeoutExceededException 的实例。

重试失败任务

要查看已插入 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

queue:flush 命令将从队列中删除所有失败的作业记录,无论失败的作业有多旧。你可以使用 --hours 选项仅删除特定小时数或更早之前失败的作业:

shell
php artisan queue:flush --hours=48

忽略缺失模型

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

你可以将闭包传递给 assertPushedassertNotPushedassertClosurePushedassertClosureNotPushed 方法,以断言已推送的作业通过了给定的「真实性测试」。如果至少推送了一项通过给定真值测试的作业,则断言将成功:

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

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

    // Perform order shipping...

    // Assert a job was pushed twice...
    Queue::assertPushedTimes(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::assertPushedTimes(ShipOrder::class, 2);
}

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

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

测试任务链

要测试作业链,你需要利用 Bus 外观的伪造功能。 Bus 门面的 assertChained 方法可用于断言 chain of jobs 已调度。 assertChained 方法接受一系列链接作业作为其第一个参数:

php
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 将确保作业实例属于同一类,并且与应用程序分派的链接作业​​具有相同的属性值:

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

你可以使用 assertDispatchedWithoutChain 方法来断言作业是在没有作业链的情况下推送的:

php
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 定义来断言链式批次符合你的期望:

php
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 的实例,该实例可用于检查批处理中的作业:

php
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 方法可用于待处理的批次,以验证该批次是否包含预期作业。该方法接受作业实例、类名或闭包的数组:

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

当使用闭包时,闭包将接收作业实例。预期的作业类型将从闭包的类型提示中推断出来:

php
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 方法来断言已调度给定数量的批次:

php
Bus::assertBatchCount(3);

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

php
Bus::assertNothingBatched();

测试任务 / 批次交互

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

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

$job->handle();

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

测试任务 / 队列交互

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

一旦作业的队列交互被伪造,你就可以在作业上调用 handle 方法。调用作业后,可以使用各种断言方法来验证作业的队列交互:

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 方法,你可以指定在工作线程尝试从队列中获取作业之前执行的回调。例如,你可以注册一个闭包来回滚之前失败的作业留下的任何事务:

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

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