本文へ移動
Laravel Tips

キュー(時間のかかる処理をあとで動かす)

時間のかかる仕事をあとまわしにして裏で動かす「キュー」の使い方を、ジョブの作り方・送り方・ワーカーの動かし方・失敗したときの扱い・テストまで説明します。

キューは、時間のかかる仕事を「順番待ちの列」に並べておき、裏で少しずつ片づけてもらうしくみです。お店で言えば、注文を受けたらすぐ番号札を渡し、料理は厨房で順に作るようなものです。たとえば、アップロードされた CSV ファイルを読み取って保存する仕事は、ふつうのリクエストの中でやると時間がかかります。キューに任せれば、画面はすぐ返せます。列に並べる1つ1つの仕事をジョブ、列から仕事を取り出して動かすプログラムをワーカーと呼びます。

キューのしくみと設定#

Laravel のキューは、Amazon SQS・Redis・ふつうのデータベースなど、いろいろな保存先(バックエンド)を同じ書き方で使えます。設定は config/queue.php にあります。

この設定には、保存先ごとの設定(ドライバー)が入っています。標準で使えるのは次のとおりです。

ドライバー 説明
database データベースの表に仕事を入れる
Amazon SQS Amazon のキューサービス
Redis Redis というデータの置き場を使う
Beanstalkd Beanstalkd というキューの専用ソフトを使う
sync 並べずにその場ですぐ実行する(開発やテスト向け)
null 並べた仕事を捨てる

補足

Laravel Horizon は、Redis で動くキューを見張ったり設定したりするための画面つきの道具です。

接続とキューのちがい#

config/queue.php の connections は、「どの保存先につなぐか」の設定です。1つの接続の中には、複数の「キュー」(仕事の山)を作れます。各接続の設定には queue という項目があり、送り先を決めなかったときに使う山の名前になります。

php
use App\Jobs\ProcessPodcast;

// 既定の接続の、既定のキューに送る
ProcessPodcast::dispatch();

// 既定の接続の "emails" キューに送る
ProcessPodcast::dispatch()->onQueue('emails');

山を分けると、急ぎの仕事を先に処理できます。たとえば high という山を作ったら、ワーカーに「high を先に見てね」と頼めます。

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

ドライバーごとの準備#

データベース#

database ドライバーには、仕事を入れる表が要ります。新しいアプリには、たいていマイグレーション(表を作る手順書。マイグレーション)の 0001_01_01_000002_create_jobs_table.php が入っています。無いときは次のコマンドで作ります。

bash
php artisan make:queue-table

php artisan migrate

Redis#

redis ドライバーを使うには、config/database.php に Redis の接続を設定します。

注意

Redis の設定項目 serializer と compression は、キューの redis ドライバーでは使えません。

Redis Cluster(Redis を複数台で動かす形)を使うときは、キューの名前にキーのハッシュタグ({ と } で囲んだ部分)を入れます。1つのキューのデータを、同じ場所にまとめておくために必要です。

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

block_for は、仕事が来るまで何秒待つかを決める設定です。何度も問い合わせ続けるより、少し待つほうが効率がよいことがあります。次の例は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,
],

注意

block_for を 0 にすると、仕事が来るまで永久に待ちます。そのあいだ、SIGTERM などの止める合図(シグナル)も、次の仕事が終わるまで受け付けません。

SQS の大きなデータの置き場#

Amazon SQS には、1つのメッセージの大きさに上限があります。上限を超えそうなデータは、キャッシュに保存して、SQS には「置き場所の目印」だけを送れます。使うには、SQS の設定に overflow を足します。

php
'sqs' => [
    'driver' => 'sqs',
    'key' => env('AWS_ACCESS_KEY_ID'),
    'secret' => env('AWS_SECRET_ACCESS_KEY'),
    'prefix' => env('SQS_PREFIX', 'https://sqs.us-east-1.amazonaws.com/your-account-id'),
    'queue' => env('SQS_QUEUE', 'default'),
    'suffix' => env('SQS_SUFFIX'),
    'region' => env('AWS_DEFAULT_REGION', 'us-east-1'),
    'after_commit' => false,
    'overflow' => [
        'enabled' => env('SQS_OVERFLOW_ENABLED', false),
        'store' => env('SQS_OVERFLOW_STORE'),
        'always' => false,
        'delete_after_processing' => true,
        'flush_on_clear' => env('SQS_OVERFLOW_FLUSH_ON_CLEAR', false),
    ],
],
項目 説明
enabled この機能を使うか
store データをしまうキャッシュの名前
always true なら、大きさに関係なく全部をキャッシュにしまう
delete_after_processing 仕事が無事に終わったら、しまったデータを消す(既定の動き)
flush_on_clear queue:clear で SQS を空にするとき、キャッシュも空にする

大きさが 1 MB 以上のデータが、キャッシュにしまわれます。ワーカーが処理するまでデータが残るキャッシュを選んでください。flush_on_clear は、キャッシュの中身を全部消すことがあるので、使うなら専用のキャッシュを用意します。

そのほかのドライバーに必要なもの#

次のものを Composer(PHP の部品を入れる道具)で入れます。

ドライバー 入れるもの
Amazon SQS aws/aws-sdk-php ~3.0
Beanstalkd pda/pheanstalk ^7.0|^8.0
Redis predis/predis ~3.0、または phpredis という PHP の拡張
MongoDB mongodb/laravel-mongodb

ジョブを作る#

ジョブのクラスを作る#

ジョブのクラスは、ふつう app/Jobs に置きます。フォルダが無ければ、次のコマンドで作られます。

bash
php artisan make:job ProcessPodcast

できたクラスは ShouldQueue というインターフェイス(決まりごと)を使っていて、「これは列に並べて、あとで動かす仕事だ」と Laravel に伝えます。

補足

ジョブのひな形(スタブ)は、書き換えて自分用にできます。

クラスの形#

ジョブには、ふつう 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 のモデル(データベースの1行を表すクラス)をそのまま渡せます。Queueable を使っているので、モデルと、読み込み済みのリレーション(表どうしのつながり)は、列に入れるときと取り出すときに自動で変換されます。列に入るのはモデルの ID だけです。仕事を動かすときに、データベースからモデルを取り直します。そのため、列に入るデータが小さくて済みます。

handle に必要な部品を渡してもらう#

handle の引数に型を書くと、サービスコンテナ(クラスを作って渡してくれる道具箱)が、必要な部品を自動で渡してくれます。

渡し方を自分で決めたいときは、コンテナの bindMethod を使います。ふつうは 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));
});

注意

画像の中身のような生のバイナリデータは、ジョブに渡す前に base64_encode で文字に変えてください。そのままだと、列に入れるときの JSON への変換に失敗することがあります。

モデルのリレーションの扱い#

読み込み済みのリレーションも、ジョブに入るときに一緒に変換されます。そのため、データが大きくなることがあります。また、取り出すときは、リレーションが全部取り直されます。入れる前に付けていた絞り込みの条件は消えます。リレーションの一部だけを使いたいなら、ジョブの中でもう一度絞り込んでください。

リレーションを変換したくないときは、モデルに withoutRelations を呼びます。

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

特定のリレーションだけ外すなら、withoutRelation を使います。

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

コンストラクターの引数をそのままプロパティにする書き方(PHP のコンストラクタープロモーション)なら、PHP の属性 WithoutRelations を付けます。

php
use Illuminate\Queue\Attributes\WithoutRelations;

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

クラス全体に付ければ、すべてのモデルのリレーションが外れます。

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,
    ) {}
}

モデルのコレクション(入れ物)や配列を受け取るジョブでは、取り出すときにリレーションが元に戻りません。大量のモデルを扱うときに、資源を使いすぎないためです。

同じジョブを1つだけにする(ユニークなジョブ)#

注意

ユニークなジョブには、ロック(同時に1つだけが使える鍵)に対応したキャッシュが要ります。対応しているのは memcached・redis・dynamodb・database・file・array です。

注意

ユニークの決まりは、バッチ(あとで説明する、ジョブをまとめて送るしくみ)の中のジョブには効きません。

同じジョブが、列に同時に1つしか入らないようにできます。ジョブのクラスに ShouldBeUnique を付けます。追加で書くメソッドはありません。

php
<?php

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

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

この例では、同じジョブがすでに列にあって終わっていなければ、新しく送っても入りません。

「何をもって同じとするか」の目印(キー)や、ユニークでいる時間の上限も決められます。PHP の属性 UniqueFor と、uniqueId メソッドを使います。

php
<?php

namespace App\Jobs;

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

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

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

この例では、商品の ID が同じジョブは、先のジョブが終わるまで無視されます。さらに、1時間たっても終わっていなければ、ロックが外れ、同じ ID の新しいジョブを送れるようになります。

注意

複数のサーバーやコンテナからジョブを送るときは、すべてが同じキャッシュのサーバーにつながるようにしてください。そうしないと、同じかどうかを正しく判断できません。

処理が始まるまでユニークにする#

ふつう、ユニークのロックは、ジョブが終わったとき(または、再挑戦を使い切って失敗したとき)に外れます。「処理が始まる直前に外したい」ときは、ShouldBeUnique の代わりに ShouldBeUniqueUntilProcessing を使います。

php
<?php

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

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

ユニークのロックに使うキャッシュ#

ShouldBeUnique のジョブを送るとき、Laravel は uniqueId の値でロックを取ろうとします。すでに取られていれば、送りません。ロックには、ふつう既定のキャッシュを使います。別のキャッシュを使いたいときは、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');
    }
}

補足

「同時に動くのを1つに制限したいだけ」なら、あとで説明する WithoutOverlapping のジョブミドルウェアを使ってください。

連続して送られたら最後の1つだけ動かす(デバウンス)#

同じジョブが短い間に何度も送られたとき、いちばん新しい1つだけを動かせます。ジョブに PHP の属性 DebounceFor を付けます。

php
<?php

namespace App\Jobs;

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

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

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

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

この例では、同じ商品のジョブが 30 秒以内に何度送られても、最後の1つだけが動きます。

何度も送り直されて、延々と後回しになるのを防ぎたいときは、maxWait で待つ時間の上限を決めます。

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

    // ...
}

デバウンスの記録に使うキャッシュは、debounceVia メソッドで選べます。

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

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

新しいジョブに置きかえられた古いジョブは、列から取り除かれ、Illuminate\Queue\Events\JobDebounced というイベント(「〜が起きた」という知らせ)が発生します。

注意

デバウンスとユニークは、同時に使えません。DebounceFor を付けたジョブは、ShouldBeUnique を付けないでください。

注意

複数のサーバーやコンテナからデバウンスのジョブを送るときは、すべてが同じキャッシュのサーバーにつながるようにしてください。

ジョブを暗号化する#

ジョブの中身を、他人に読まれたり書き換えられたりしないようにできます。ジョブに ShouldBeEncrypted を付けると、列に入れる前に自動で暗号化されます。

php
<?php

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

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

ジョブミドルウェア#

ジョブミドルウェアは、ジョブの実行の前後に決まった処理をはさむしくみです。handle の中身をすっきりさせられます。

たとえば、Redis で「5秒に1つだけ処理する」ようにする処理を、handle に直接書くとこうなります。

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 が読みにくくなります。ほかのジョブにも同じ処理を書くことになります。そこで、この処理をジョブミドルウェアに切り出します。

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

ルートのミドルウェアと同じで、処理中のジョブと、続きを進めるための $next を受け取ります。

ジョブミドルウェアのクラスは、make:job-middleware コマンドで作れます。ジョブに付けるには、ジョブの middleware メソッドで返します。make:job で作ったジョブにはこのメソッドが無いので、自分で足します。

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

補足

ジョブミドルウェアは、キューに入れるイベントのリスナー(イベント)、メール(メール)、通知(通知)にも付けられます。

Laravel が用意しているジョブミドルウェアは次のとおりです。

ミドルウェア 説明
RateLimited 回数の制限をかける
RateLimitedWithRedis 回数の制限(Redis 向け)
WithoutOverlapping 同じ目印のジョブが重ならないようにする
ThrottlesExceptions 例外が続いたら、しばらく待たせる
ThrottlesExceptionsWithRedis 例外の制限(Redis 向け)
Release 条件に合えば、実行せずに列へ戻す
Skip 条件に合えば、そのジョブを飛ばして消す
FailOnException 特定の例外が出たら、すぐ失敗にする
SkipIfBatchCancelled バッチが取り消されていたら、実行しない

回数の制限(RateLimited)#

自分で書かなくても、回数を制限するミドルウェアが用意されています。ルートのレート制限と同じく、RateLimiter の for メソッドで制限を定義します。

たとえば、「ふつうの利用者は、データのバックアップを1時間に1回まで。特別な利用者は制限なし」にするには、AppServiceProvider の boot メソッドに書きます。

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

1分あたりの制限なら 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)は増えます。ジョブの Tries や MaxExceptions の値を、それに合わせて調整してください。あるいは、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 を使っているなら、Redis 向けに調整された、より効率のよい Illuminate\Queue\Middleware\RateLimitedWithRedis が使えます。

php
use Illuminate\Queue\Middleware\RateLimitedWithRedis;

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

使う Redis の接続は、connection メソッドで選べます。

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

ジョブが重ならないようにする(WithoutOverlapping)#

Illuminate\Queue\Middleware\WithoutOverlapping は、好きな目印(キー)ごとに、同じジョブが重なって動かないようにします。1つのデータを、一度に1つのジョブだけが書き換えるようにしたいときに役立ちます。

たとえば、利用者の信用スコアを更新するジョブで、同じ利用者の更新が重ならないようにするには、こう書きます。

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

重なったジョブを列に戻したときも、試した回数は増えます。Tries を調整してください。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()];
}

このミドルウェアは、ロックのしくみで動いています。ジョブが思わぬ形で失敗したり時間切れになったりすると、ロックが外れないことがあります。そのため、expireAfter でロックの期限を決められます。次の例は、処理が始まってから3分でロックが外れます。

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

注意

WithoutOverlapping には、ロックに対応したキャッシュが要ります。対応しているのは memcached・redis・dynamodb・database・file・array です。

ジョブの種類をまたいで同じ目印を使う#

ふつう、WithoutOverlapping が防ぐのは、同じクラスのジョブどうしの重なりです。別のクラスが同じ目印を使っても、重なりは防がれません。クラスをまたいで目印を共有したいときは、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(),
        ];
    }
}

例外が続いたら待たせる(ThrottlesExceptions)#

Illuminate\Queue\Middleware\ThrottlesExceptions は、ジョブが決まった回数の例外(処理の途中で起きたエラーの知らせ)を出したら、次の試しを一定の時間だけ遅らせます。調子の悪い外部のサービスとやりとりするジョブに向いています。

ふつうは、「いつまで試すか」を決める retryUntil(時間ベースの試し)と組み合わせます。

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

1つ目の引数は「何回の例外で待たせるか」、2つ目は「待たせたあと、何秒たってから試すか」です。上の例では、連続で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)];
}

backoff には、出た例外を受け取るクロージャ(名前のない関数)も渡せます。例外によって待ち時間を変えられます。

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

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

このミドルウェアは、内部でキャッシュを使い、ジョブのクラス名を目印にしています。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 に条件を渡すと、条件が 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 を使います。クロージャを渡せば、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 を使っているなら、Redis 向けの Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis が、より効率よく動きます。

php
use Illuminate\Queue\Middleware\ThrottlesExceptionsWithRedis;

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

接続は connection メソッドで選べます。

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

実行せずに列へ戻す(Release)#

Release ミドルウェアは、ジョブを実行せずに列へ戻します。Release::when は条件が true のとき、Release::unless は条件が false のときに戻します。

php
use Illuminate\Queue\Middleware\Release;

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

戻したときも、試した回数は増えます。Tries や MaxExceptions を調整してください。

複雑な条件は、クロージャでも渡せます。

php
use Illuminate\Queue\Middleware\Release;

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

ジョブを飛ばす(Skip)#

Skip ミドルウェアは、ジョブの中身を変えずに、ジョブを飛ばして消します。Skip::when は条件が true のとき、Skip::unless は条件が false のときに消します。

php
use Illuminate\Queue\Middleware\Skip;

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

こちらも、クロージャを渡せます。

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

条件つきで送るには、dispatchIf と dispatchUnless を使います。

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

注意

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

応答を返したあとに動かす#

deferred 接続に送ると、今のプロセス(動いているプログラム)の中で動かせます。ただし、動くのは応答をブラウザへ返したあとです。利用者を待たせずに、その場で処理できます。

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

deferred 接続は、あとで説明するフェイルオーバーのキューの既定の行き先でもあります。

background 接続も、応答を返したあとに動かします。こちらは、別に起動した PHP のプロセスで動かします。そのため、PHP-FPM(Web サーバーの裏で PHP を動かすしくみ)などのアプリ側のワーカーは、すぐ次のリクエストを受けられます。

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

まとめて送る(バルク)#

たくさんのジョブを、まとめて送れます。バッチの記録やコールバック(終わったときに呼ばれる処理)が要らず、1つ1つが独立したジョブなら、Bus ファサード(クラス名と :: で機能を呼べる窓口)の bulk を使います。Laravel が接続とキューの名前ごとに分けて、まとめて列に入れます。

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

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

送る前に準備する#

ジョブを列に入れる前に、状態を調べたいときは、Illuminate\Contracts\Queue\PreparesForDispatch を付けて prepareForDispatch メソッドを書きます。Laravel は、送る前にこのメソッドを呼びます。false を返すと、そのジョブは送られません。

php
<?php

namespace App\Jobs;

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

class SyncPodcasts implements PreparesForDispatch, ShouldQueue
{
    use Queueable;

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

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

データベースのトランザクションとジョブ#

トランザクション(複数の書き込みを、まとめて確定するか取り消す仕組み)の中でジョブを送ることもできます。ただし、注意が要ります。ワーカーが、トランザクションの確定より先にジョブを処理することがあるからです。そのときは、トランザクションの中で変えたデータがまだ反映されていなかったり、作ったはずのデータがまだ無かったりします。

対策の1つ目は、キュー接続の設定で after_commit を true にすることです。

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

true のときは、開いているトランザクションが確定するまで、Laravel がジョブを実際には送りません。トランザクションが開いていなければ、すぐ送ります。トランザクションの途中で例外が起きて取り消されたら、その中で送ったジョブは捨てられます。

補足

after_commit を true にすると、キューに入れるイベントのリスナー・メール・通知・ブロードキャストのイベントも、トランザクションの確定後に送られます。

ジョブごとに決める#

after_commit を設定していなくても、ジョブごとに afterCommit をつなげれば、確定後に送れます。

php
use App\Jobs\ProcessPodcast;

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

逆に、after_commit が true のときでも、待たずにすぐ送りたいジョブには beforeCommit を使います。

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

ジョブをつなぐ(チェーン)#

チェーンは、「最初のジョブが成功したら次、その次…」と、順番に動かすしくみです。途中で1つ失敗すると、残りは動きません。Bus ファサードの chain を使います。

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

注意

ジョブの中で $this->delete() を呼んで消しても、後ろのジョブは止まりません。チェーンが止まるのは、チェーンの中のジョブが失敗したときだけです。

チェーンの接続とキュー#

つなげたジョブの接続とキューは、onConnection と onQueue で決めます。ジョブ自身が別の接続やキューを決めているときは、そちらが優先されます。

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

チェーンにジョブを足す#

チェーンの中のジョブから、そのチェーンの先頭や最後にジョブを足せます。prependToChain と appendToChain を使います。

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

注意

チェーンのコールバックは、列に入れるために変換され、あとで動きます。そのため、中で $this は使えません。

キューと接続を選ぶ#

送り先のキュー#

ジョブを別のキュー(山)に送ると、種類ごとに分けたり、ワーカーの割り当てで優先度をつけたりできます。これは、設定ファイルの別の「接続」に送るのではなく、1つの接続の中の別のキューに送ることです。送るときに 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');
    }
}

送り先の接続#

複数のキュー接続を使うなら、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');
    }
}

onConnection と onQueue は、つなげて使えます。

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

こちらも、コンストラクターの中で決められます。

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

キューのルーティング#

Queue ファサードの route メソッドで、特定のジョブのクラスの既定の接続とキューを決められます。ジョブごとに指定を書かなくても、いつも決まったキューに入ります。

クラス名のほかに、インターフェイス・トレイト(いくつものクラスに同じ機能を足す部品)・親クラスも渡せます。そのときは、それを使っている(または継いでいる)ジョブが、すべて同じ行き先になります。ふつうは、サービスプロバイダの boot メソッドに書きます。

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

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

接続だけを決めたときは、その接続の既定のキューに入ります。

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

配列で、まとめて決めることもできます。

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

補足

キューのルーティングは、ジョブ自身が決めた行き先で上書きできます。

forward を使うと、あるキューに送られたジョブを、別のキューや接続へ転送できます。ジョブや送る側のコードを変えずに、キューの置き場所を変えたいときに便利です。

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

配列で、まとめて転送もできます。

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

ジョブ自身が決めた接続は、転送先の接続より優先されます。

試す回数と制限時間#

試す回数#

「試す回数」(attempts)は、キューの大事な考え方です。ジョブを送ると列に入り、ワーカーが取り出して実行しようとします。これが1回の試しです。

ただし、「試した」からといって、handle が動いたとは限りません。次のどれでも、1回分が使われます。

  • ジョブの実行中に、処理されなかった例外が出た
  • ジョブの中で $this->release() を呼んで、列へ戻した
  • WithoutOverlapping や RateLimited などのミドルウェアがロックを取れず、列へ戻した
  • ジョブが時間切れになった
  • handle が例外なしで最後まで動いた

同じジョブを延々と試し続けたくはないので、Laravel は回数や時間の上限を決める方法を用意しています。

補足

既定では、ジョブは1回しか試されません。WithoutOverlapping や RateLimited を使ったり、手動で列へ戻したりするなら、たいてい tries で回数を増やすことになります。

1つ目の方法は、queue:work の --tries です。ジョブ自身が回数を決めていなければ、ワーカーが処理するすべてのジョブに効きます。

bash
php artisan queue:work --tries=3

試す回数を使い切ったジョブは、「失敗したジョブ」として扱われます。--tries=0 にすると、無制限に試します。

ジョブのクラスに、PHP の属性 Tries で回数を決めるやり方もあります。ジョブに決めた値が、コマンドの --tries より優先されます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Tries;

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

回数を、その場の状況で決めたいときは、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);
}

retryUntil と tries の両方があるときは、retryUntil が優先されます。

補足

Tries の属性や retryUntil メソッドは、キューに入れるイベントのリスナーや通知にも書けます。

例外の回数の上限#

「何度でも試してよいが、処理されなかった例外が決まった回数に達したら失敗にする」こともできます(release で戻しただけの再挑戦は、数えません)。ジョブに、PHP の属性 Tries と MaxExceptions を付けます。

php
<?php

namespace App\Jobs;

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

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

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

この例では、Redis のロックが取れなければ10秒おいて、最大25回まで試します。ただし、処理されなかった例外が3回出たら失敗になります。

ワーカーのプロセスが落ちたり(メモリ不足で止められたときなど)、強制終了されたりして終わった試しは、既定では、例外の回数に数えません。数えたいときは、PHP の属性 CountCrashesAsExceptions を付けます。

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

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

この属性があると、ワーカーはジョブの処理中、アプリのキャッシュに目印を置きます。次の試しのときに目印が残っていれば、前の試しを例外として数えます。

例外によって再挑戦をやめる#

ある例外が出たら、再挑戦せずにすぐ失敗にしたいことがあります。bootstrap/app.php の dontRetry で、その例外の種類を並べます。

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

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

もっと細かく決めたいときは、dontRetryWhen にクロージャを渡します。true を返すと、そのジョブは失敗とされ、再挑戦されません。

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

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

制限時間(タイムアウト)#

ジョブにかかる時間は、だいたい分かっていることが多いはずです。Laravel では、制限時間(タイムアウト)を決められます。既定は60秒です。これより長く処理していると、そのジョブを動かしているワーカーは、エラーで終了します。ふつうは、サーバーに設定したプロセスの管理役が、ワーカーを自動で起動し直します。

制限時間は、コマンドの --timeout で決められます。

bash
php artisan queue:work --timeout=30

制限時間切れを続けて、試す回数を使い切ると、そのジョブは失敗になります。

ジョブのクラスに、PHP の属性 Timeout で決めることもできます。ジョブに決めた値が、コマンドの値より優先されます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Timeout;

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

ソケットや外へ出る HTTP 通信など、待ちが入る処理は、決めた制限時間を守らないことがあります。そういう処理を使うときは、その道具の側でも必ず制限時間を決めてください。たとえば Guzzle(HTTP の部品)なら、接続とリクエストの制限時間を決めます。

注意

ジョブの制限時間を使うには、PHP の PCNTL という拡張が要ります。また、ジョブの制限時間は、「再挑戦までの時間」(後で説明する retry_after)より短くしてください。そうしないと、終わる前に再挑戦されることがあります。queue:work に --once をつけたときは、--timeout は効きません。

時間切れで失敗にする#

時間切れになったジョブを、失敗として扱いたいときは、PHP の属性 FailOnTimeout を付けます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\FailOnTimeout;

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

補足

ふつう、時間切れのジョブは1回分を使い、再挑戦が許されていれば列に戻ります。FailOnTimeout を付けると、tries がいくつでも、再挑戦されません。

SQS の FIFO キューと公平なキュー#

Laravel は、Amazon SQS の FIFO(先に入れたものから出る)キューと、公平(fair)なキューに対応しています。FIFO キューは、送った順に処理し、同じ仕事が二重に処理されないようにします。

FIFO キューには、メッセージグループの ID が要ります。同じ ID のジョブは順番に処理され、ちがう ID のジョブは同時に処理できます。送るときに onGroup で ID を決めます。

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

グループを決めずに FIFO キューへ送ると、キューの名前がグループの ID になります。

二重に処理されないようにするための「重複の目印」も決められます。ジョブに deduplicationId メソッドを書きます。

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

公平なキュー#

SQS の標準のキューでメッセージグループを決めると、公平なキューとして動きます。グループを割り当てれば、利用者や仕事の種類のあいだで、公平に配られます。Laravel に、追加の設定は要りません。

送るときの onGroup のかわりに、ジョブに messageGroup メソッドを書いても決められます。

php
<?php

namespace App\Jobs;

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

class ProcessOrder implements ShouldQueue
{
    use Queueable;

    // ...

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

FIFO とリスナー・メール・通知#

FIFO キューを使うときは、リスナー・メール・通知にも、メッセージグループを決める必要があります。かわりに、FIFO ではないキューに送ってもかまいません。

キューに入れるイベントのリスナーには、messageGroup メソッドを書きます。重複の目印をつくる deduplicator メソッドも、必要なら書けます。イベントを受け取り、目印を作るクロージャを返します。

php
<?php

namespace App\Listeners;

use App\Events\OrderShipped;
use Closure;

class SendShipmentNotification
{
    // ...

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

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

FIFO キューへ送るメールは、onGroup と、必要なら withDeduplicator を呼んでから送ります。

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 キューへ送る通知も、同じです。

php
use App\Notifications\InvoicePaid;

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

$user->notify($invoicePaid);

キューのフェイルオーバー#

failover ドライバーは、ジョブを列に送るのに失敗したとき、自動で次の接続に送り直してくれます。本番で、キューが止まると困るときに役立ちます。

設定では、failover ドライバーと、試す接続の名前を順番に並べます。新しいアプリの config/queue.php には、例が最初から入っています。

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

使うには、.env で、この接続を既定のキュー接続にします。

ini
QUEUE_CONNECTION=failover

次に、並べた接続ごとに、ワーカーを1つ以上起動します。

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

補足

sync・background・deferred のドライバーの接続は、今の PHP のプロセスの中で処理するので、ワーカーは要りません。

接続の操作が失敗してフェイルオーバーが動くと、Illuminate\Queue\Events\QueueFailedOver というイベントが発生します。キューの接続の失敗を、記録や通知に使えます。

補足

Laravel Horizon が管理するのは、Redis のキューだけです。フェイルオーバーの一覧に database があるなら、Horizon とは別に、ふつうの php artisan queue:work database も動かしてください。

エラーの扱い#

ジョブの処理中に例外が出ると、ジョブは自動で列に戻され、もう一度試されます。アプリで決めた最大の試す回数に達するまで続きます。最大の回数は、queue:work の --tries か、ジョブのクラスで決めます。

手動で列に戻す#

あとでもう一度試したいときは、release を呼んで、自分で列に戻せます。

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

    $this->release();
}

既定では、すぐ処理できる状態で戻ります。秒数や日時を渡せば、その時間がたつまで待たせられます。

php
$this->release(10);

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

手動で失敗にする#

ジョブを「失敗」にしたいときは、fail を呼びます。

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

    $this->fail();
}

捕まえた例外が原因なら、その例外を渡せます。文字列を渡すと、例外に変えてくれます。

php
$this->fail($exception);

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

補足

失敗したジョブの扱いは、後の「失敗したジョブの扱い」を見てください。

特定の例外で失敗にする#

FailOnException ジョブミドルウェアを使うと、特定の例外が出たときに、再挑戦をやめてすぐ失敗にできます。外部の 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\Attributes\Tries;
use Illuminate\Queue\Middleware\FailOnException;
use Illuminate\Support\Facades\Http;

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

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

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

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

        // ...
    }

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

ジョブのバッチ#

バッチは、ジョブのグループを並行して動かし、全部が終わったら何かをするしくみです。

使う前に、バッチの記録(完了した割合など)を入れる表のマイグレーションを作ります。

bash
php artisan make:queue-batches-table

php artisan migrate

バッチにできるジョブを作る#

ふつうのジョブとして作り、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 で送ります。バッチは、終わったときのコールバックと組み合わせて使うのが基本です。コールバックは、次のメソッドで決めます。どれも Illuminate\Bus\Batch を受け取ります。

メソッド 動くとき
before バッチが作られたが、まだジョブが追加されていないとき
progress 1つのジョブが成功したとき
then 全部のジョブが成功したとき
catch バッチの中でジョブの失敗が見つかったとき
finally バッチの実行が終わったとき

ワーカーを複数動かすと、バッチのジョブは並行して処理されます。そのため、終わる順番は、入れた順番と同じとは限りません。順番に動かしたいときは、後の「チェーンとバッチ」を見てください。

次は、CSV ファイルの決まった行数ずつを処理する、5つのジョブのバッチです。

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;

$batch->id で取り出せるバッチの ID は、送ったあとで、バッチの様子を調べるときに使います。

注意

バッチのコールバックは、変換されてあとで動くので、中で $this は使えません。また、バッチのジョブはデータベースのトランザクションの中で動きます。そのため、動かしただけで勝手に確定してしまう SQL を、ジョブの中で動かさないでください。

バッチに名前をつける#

バッチに名前をつけると、Laravel Horizon や Laravel Telescope などの道具が、分かりやすい情報を出してくれます。バッチを定義するとき、name を呼びます。

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

バッチの接続とキュー#

バッチのジョブの接続とキューは、onConnection と onQueue で決めます。バッチのジョブは、すべて同じ接続とキューで動かす必要があります。

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

チェーンとバッチ#

バッチの中に、配列でチェーンを入れられます。たとえば、2本のチェーンを並行して動かし、両方が終わったらコールバックを動かせます。

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

逆に、チェーンの中にバッチを入れることもできます。次の例は、まず複数のポッドキャストを公開するバッチを動かし、そのあとで、お知らせを送るバッチを動かします。

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

バッチにジョブを足す#

バッチの中のジョブから、同じバッチにジョブを足せます。数千のジョブを送るとき、リクエストの中では時間がかかりすぎます。そこで、まず「ジョブを足す役」のジョブだけのバッチを送り、その中で、残りのジョブを足していきます。

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

LoadImportBatch の中で、batch メソッドから取り出したバッチに add します。

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

注意

バッチにジョブを足せるのは、そのバッチに属するジョブの中だけです。

バッチの様子を調べる#

バッチのコールバックに渡される 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();
名前 説明
$batch->id バッチの UUID(ほかと重ならない ID)
$batch->name バッチの名前(つけたとき)
$batch->totalJobs バッチに入れたジョブの数
$batch->pendingJobs まだ処理されていないジョブの数
$batch->failedJobs 失敗したジョブの数
processedJobs() ここまでに処理したジョブの数
progress() 終わった割合(0から100)
finished() 実行が終わったか
cancel() バッチの実行を取り消す
cancelled() バッチが取り消されたか

ルートからバッチを返す#

Batch は JSON に変換できるので、ルートからそのまま返せば、進み具合などの情報を JSON で返せます。画面に進み具合を出したいときに便利です。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);
});

バッチを取り消す#

バッチの実行を取り消したいときは、Batch の cancel を呼びます。

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

        return;
    }

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

上の例のように、バッチのジョブは、続ける前に、バッチが取り消されていないかを確かめるのがふつうです。ジョブに SkipIfBatchCancelled のミドルウェアを付ければ、取り消されたバッチのジョブを、自動で処理せずに済みます。

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

失敗したジョブをやり直す#

queue:retry-batch コマンドで、そのバッチの失敗したジョブを、まとめてやり直せます。バッチの UUID を渡します。

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

古いバッチを消す#

何もしないと、job_batches の表の記録は、どんどん増えます。毎日1回、queue:prune-batches コマンドが動くように、スケジュールしてください。

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

ジョブが失敗して、やり直しも成功していないなど、成功せずに終わらなかったバッチの記録も、表に残ります。unfinished オプションで、それも消せます。

php
use Illuminate\Support\Facades\Schedule;

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

取り消されたバッチの記録も、cancelled オプションで消せます。

php
use Illuminate\Support\Facades\Schedule;

Schedule::command('queue:prune-batches --hours=48 --cancelled=72')->daily();
オプション 説明
--hours 終わってから何時間たったバッチを消すか
--unfinished 終わっていないバッチを、何時間たったら消すか
--cancelled 取り消されたバッチを、何時間たったら消すか

バッチの記録を DynamoDB に置く#

バッチの記録は、データベースのかわりに DynamoDB(Amazon のデータベース)にも置けます。ただし、記録を入れるための DynamoDB の表は、自分で作る必要があります。

表の名前は、ふつう job_batches です。アプリの queue 設定の queue.batching.table の値にそろえて名づけます。

DynamoDB の表の設定#

job_batches の表には、文字列の主パーティションキー application と、文字列の主ソートキー id が要ります(どちらも、DynamoDB で1件ずつを見分けるためのキー)。application には、アプリの app 設定の name の値が入ります。アプリ名がキーに入るので、複数の Laravel アプリで同じ表を使えます。

自動で消す機能を使うなら、表に ttl の項目も定義してください。

DynamoDB の設定#

Laravel から Amazon DynamoDB を使うために、AWS SDK を入れます。

bash
composer require aws/aws-sdk-php

そして、queue.batching.driver を dynamodb にします。batching の設定に key・secret・region も書きます。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 の TTL(期限が来たら自動で消す機能)を使います。

表に ttl の項目があるなら、消し方を設定に書きます。queue.batching.ttl_attribute は、TTL を入れる項目の名前です。queue.batching.ttl は、最後に更新されてから何秒で、記録を消してよいかを決めます。

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',
    '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();
});

name で名前をつけると、キューの様子を見る画面や、queue:work の出力に、その名前が出ます。

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

catch でクロージャを渡すと、設定した再挑戦を使い切っても成功しなかったときに、それが動きます。

php
use Throwable;

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

注意

catch のクロージャは、変換されてあとで動くので、中で $this は使えません。

ワーカーを動かす#

queue:work コマンド#

queue:work コマンドで、ワーカーを起動します。ワーカーは、新しいジョブが列に入るたびに処理します。いったん起動すると、自分で止めるか、ターミナルを閉じるまで動き続けます。

bash
php artisan queue:work

補足

queue:work を裏でずっと動かし続けるには、Supervisor のようなプロセスの管理役を使って、止まらないようにします。

-v をつけると、処理したジョブの ID・接続の名前・キューの名前が出力に出ます。

bash
php artisan queue:work -v

ワーカーは、長く動き続けるプロセスで、起動したときのアプリの状態をメモリに持ちつづけます。そのため、起動したあとにコードを変えても、気づきません。デプロイのときは、ワーカーを再起動してください。また、アプリが作ったり変えたりした static(クラスに持たせた共有の値)の状態は、ジョブのあいだで自動では元に戻りません。

かわりに queue:listen を使うと、コードを変えても、ワーカーを再起動する必要はありません。ただし、queue:work よりずっと効率が悪くなります。

bash
php artisan queue:listen

ワーカーを複数動かす#

1つのキューに複数のワーカーを割り当てて、同時に処理するには、queue:work を複数起動します。手元では、ターミナルのタブを複数開きます。本番では、プロセスの管理役の設定で決めます。Supervisor なら、numprocs で決められます。

接続とキューを選ぶ#

ワーカーが使うキュー接続も選べます。コマンドに渡す接続の名前は、config/queue.php の接続のどれかにそろえます。

bash
php artisan queue:work redis

既定では、queue:work は、その接続の既定のキューのジョブだけを処理します。特定のキューだけを処理することもできます。たとえば、redis 接続の emails キューだけなら、次のようにします。

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

処理する数を決める#

--once は、列から1つのジョブだけ処理して終わります。

bash
php artisan queue:work --once

--max-jobs は、決めた数のジョブを処理して終わります。Supervisor と組み合わせると、一定の数ごとにワーカーが自動で起動し直され、たまったメモリを手放せます。

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

列が空になるまで処理して終わる#

--stop-when-empty は、列のジョブを全部処理してから、きれいに終わります。Docker のコンテナで、列が空になったらコンテナも止めたいときに便利です。

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

決めた秒数だけ処理する#

--max-time は、決めた秒数のあいだ処理してから終わります。Supervisor と組み合わせると、一定の時間ごとに起動し直され、たまったメモリを手放せます。

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

ワーカーの休む時間#

列にジョブがあるあいだは、ワーカーは休まず処理を続けます。ジョブが無いときは、sleep で決めた秒数だけ休みます。休んでいるあいだは、新しいジョブを処理しません。

bash
php artisan queue:work --sleep=3

メンテナンスモードとキュー#

アプリがメンテナンスモードのあいだは、列のジョブは処理されません。メンテナンスモードが終われば、また処理されます。

メンテナンスモードでも処理させたいときは、--force を使います。

bash
php artisan queue:work --force

資源の扱い#

ワーカーは、ジョブの前に、フレームワークを起動し直しません。そのため、ジョブが終わるたびに、重い資源を手放してください。たとえば、GD ライブラリで画像を加工したなら、終わったときに imagedestroy でメモリを空けます。

queue:work のオプションのまとめ#

オプション 説明
-v 処理したジョブの ID・接続・キューを出す
--queue= 処理するキューを選ぶ(カンマで区切って並べると優先順になる)
--tries= 試す回数の上限(0 なら無制限)
--timeout= ジョブの制限時間(秒)
--backoff= 失敗してから再挑戦するまでの秒数
--once 1つのジョブだけ処理して終わる
--max-jobs= 決めた数のジョブを処理して終わる
--stop-when-empty 列が空になったら終わる
--max-time= 決めた秒数のあいだ処理して終わる
--sleep= ジョブが無いときに休む秒数
--force メンテナンスモードでも処理する

キューの優先順位#

キューによって、処理する順番を変えたいことがあります。たとえば、redis 接続の既定のキューを low にしておき、急ぎのジョブだけを high に送ります。

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

high のジョブを全部処理してから low に進むワーカーを起動するには、カンマで区切ってキュー名を並べます。

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

ワーカーとデプロイ#

ワーカーは長く動き続けるので、再起動しないと、コードの変更に気づきません。いちばん簡単なデプロイの方法は、デプロイのとき、ワーカーを再起動することです。queue:restart で、全部のワーカーを、きれいに再起動させられます。

bash
php artisan queue:restart

このコマンドは、全ワーカーに「今のジョブが終わったら終了して」と伝えます。そのため、途中のジョブが失われません。ワーカーは終了するので、Supervisor のようなプロセスの管理役で、自動で起動し直すようにしておいてください。

補足

再起動の合図はキャッシュに保存されます。この機能を使う前に、キャッシュが正しく設定されていることを確かめてください。

ワーカーへの合図に反応する#

処理中のワーカーが、SIGQUIT・SIGTERM・SIGINT などの止める合図を受けても、今のジョブは最後まで終わらせてから終了します。ただ、サーバーやコンテナの管理役にプロセスを止められる前に、ジョブ側で反応したいこともあります。たとえば、長い取り込みのジョブが、新しいデータの読み込みをやめ、途中の進み具合を保存したいときです。

ジョブの中で合図に反応するには、Illuminate\Contracts\Queue\Interruptible を付けて、interrupted メソッドを書きます。受けた合図の番号が、interrupted に渡されます。

php
<?php

namespace App\Jobs;

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

class ImportProducts implements ShouldQueue, Interruptible
{
    use Queueable;

    protected bool $shouldStop = false;

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

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

            // Import the product row...
        }

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

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

interrupted が呼ばれるのは、ジョブが動いているあいだに、ワーカーが合図を受けたときだけです。制限時間や、ジョブの failed メソッドのかわりにはなりません。

ジョブの期限と制限時間#

ジョブの期限(retry_after)#

config/queue.php の各キュー接続には、retry_after があります。処理中のジョブが、戻されも消されもせずに何秒たったら、列に戻して再挑戦するかを決めます。たとえば 90 なら、90秒たってもジョブが終わらなければ、列に戻されます。ふつうは、ジョブにかかる最大の時間に合わせて決めます。

注意

retry_after が無いキュー接続は、Amazon SQS だけです。SQS では、AWS の管理画面で決める「可視性タイムアウト」(取り出したメッセージを、ほかから見えなくしておく時間)の既定値がたつと、再挑戦します。

ワーカーの制限時間#

queue:work の --timeout は、既定で60秒です。これより長く処理していると、そのジョブを動かしているワーカーは、エラーで終了します。ふつうは、プロセスの管理役が、ワーカーを自動で起動し直します。

bash
php artisan queue:work --timeout=60

retry_after と --timeout は別のものです。2つが組になって、ジョブが失われず、しかも成功が1回だけになるようにしています。

注意

--timeout は、retry_after より、少なくとも数秒短くしてください。そうすれば、止まってしまったジョブのワーカーを、再挑戦の前に、かならず終わらせられます。--timeout のほうが長いと、ジョブが二重に処理されることがあります。

ワーカーを一時停止・再開する#

システムの保守などで、ワーカーを止めずに、新しいジョブの処理だけを一時的に止めたいことがあります。queue:pause と queue:continue コマンドを使います。

特定のキューを止めるには、接続の名前とキューの名前を渡します。

bash
php artisan queue:pause database:default

この例では、database が接続の名前、default がキューの名前です。止めたキューを処理しているワーカーは、今のジョブは最後までやりますが、再開するまで、新しいジョブは取りません。

すべての接続のすべてのキューを止めるには、--all を使います。

bash
php artisan queue:pause --all

止めたキューを再開するには、queue:continue を使います。

bash
php artisan queue:continue database:default

すべての接続のすべてのキューを再開するには、queue:resume に --all をつけます。

bash
php artisan queue:resume --all

再開すると、ワーカーはすぐ、そのキューの新しいジョブを処理しはじめます。すべてを再開しても、1つずつ個別に止めたキューは再開されません。キューを止めても、ワーカーのプロセス自体は止まりません。指定したキューの新しいジョブを取らなくなるだけです。

コマンド 説明
queue:pause 接続:キュー 指定したキューの処理を止める
queue:pause --all すべてのキューの処理を止める
queue:continue 接続:キュー 止めたキューを再開する
queue:resume --all すべてのキューを再開する

再起動と一時停止の合図#

既定では、ワーカーはジョブを1つ処理するたびに、キャッシュに「再起動」「一時停止」の合図が来ていないかを調べます。queue:restart と queue:pause に応えるために大事ですが、わずかに負担になります。

これらの機能が要らず、速さを優先したいなら、Queue ファサードの withoutInterruptionPolling で、全体でこの確認をやめられます。ふつうは、AppServiceProvider の boot メソッドで呼びます。

php
use Illuminate\Support\Facades\Queue;

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

再起動と一時停止を別々にやめるには、Illuminate\Queue\Worker の static のプロパティ $restartable と $pausable を false にします。

php
use Illuminate\Queue\Worker;

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

注意

確認をやめると、ワーカーは、やめた機能に対応するコマンド(queue:restart や queue:pause)に応えなくなります。

Supervisor の設定#

本番では、queue:work を動かし続ける方法が要ります。ワーカーは、制限時間を超えたり、queue:restart を実行したりなど、いろいろな理由で止まります。

そのため、queue:work が終了したことに気づいて、自動で起動し直してくれる「プロセスの管理役」を設定します。管理役は、同時に動かす queue:work の数も決められます。Supervisor は、Linux でよく使われる管理役です。

Supervisor を入れる#

Supervisor は Linux 向けの管理役で、queue:work が止まったら、自動で起動し直します。Ubuntu では、次のコマンドで入れます。

bash
sudo apt-get install supervisor

補足

Supervisor を自分で設定して管理するのが大変なら、Laravel Cloud という、キューのワーカーを動かすことまで任せられるサービスもあります。

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 は、queue:work を8つ動かして、全部を見張り、止まったら起動し直すという指示です。command は、使いたいキュー接続やオプションに合わせて書き換えます。

注意

stopwaitsecs は、いちばん長いジョブの秒数より大きくしてください。そうしないと、ジョブが終わる前に、Supervisor がジョブを終わらせてしまうことがあります。

Supervisor を起動する#

設定ファイルを作ったら、次のコマンドで設定を読み込んで、プロセスを起動します。

bash
sudo supervisorctl reread

sudo supervisorctl update

sudo supervisorctl start "laravel-worker:*"

失敗したジョブの扱い#

ジョブは失敗することもあります。心配はいりません。Laravel には、試す回数の上限を決める方法があります。列に入れて裏で動かしたジョブが、上限の回数を超えて失敗すると、failed_jobs というデータベースの表に入ります。その場で動かしたジョブ(dispatchSync)が失敗しても、この表には入りません。そのときの例外は、アプリがすぐに扱います。

failed_jobs を作るマイグレーションは、新しいアプリには、たいてい入っています。無いときは、make:queue-failed-table で作ります。

bash
php artisan make:queue-failed-table

php artisan migrate

ワーカーを動かすときに、試す回数の上限を --tries で決められます。--tries を決めなければ、ジョブは1回だけ、またはジョブの Tries の属性で決めた回数だけ試されます。

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

--backoff を使うと、例外が出たジョブを、何秒待ってから再挑戦するか決められます。既定では、すぐ列に戻されます。

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

ジョブごとに待つ秒数を決めるなら、PHP の属性 Backoff を付けます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Backoff;

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

待つ時間を複雑に決めたいときは、backoff メソッドを書きます。

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

秒数の配列を渡すと、だんだん待ち時間を伸ばす「指数的な」待ち方にできます。次の例では、1回目の再挑戦まで1秒、2回目まで5秒、3回目まで10秒です。それ以降は、試せる回数が残っているかぎり、10秒ずつ待ちます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\Backoff;

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

失敗したあとの後始末#

ジョブが失敗したとき、利用者に知らせたり、途中までやった処理を元に戻したりしたいことがあります。ジョブのクラスに、failed メソッドを書きます。失敗の原因の Throwable が渡されます。

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

注意

failed を呼ぶ前に、ジョブの新しいインスタンス(クラスから作った実物)が作られます。そのため、handle の中でプロパティに書き込んだ変更は、失われています。

「失敗したジョブ」は、処理されなかった例外が出たものだけではありません。許された試しを使い切ったジョブも、失敗とみなされます。試しは、次のようなことで使われます。

  • ジョブが時間切れになった
  • ジョブの実行中に、処理されなかった例外が出た
  • ジョブが、手動またはミドルウェアで、列へ戻された

最後の試しが、実行中の例外で失敗したなら、その例外が failed に渡されます。最大の試す回数に達して失敗したときは、$exception は Illuminate\Queue\MaxAttemptsExceededException になります。制限時間を超えて失敗したときは、Illuminate\Queue\TimeoutExceededException になります。

失敗したジョブをやり直す#

failed_jobs に入った失敗したジョブを全部見るには、queue:failed を使います。

bash
php artisan queue:failed

一覧には、ジョブの ID・接続・キュー・失敗した時刻などが出ます。ジョブの ID は、やり直すときに使います。たとえば、ID が ce7bb17c-cdd8-41f0-a8ec-7b4fef4e5ece のジョブをやり直すには、次のようにします。

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

ID は、複数渡せます。

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

特定のキューの、失敗したジョブを全部やり直すこともできます。

bash
php artisan queue:retry --queue=name

失敗したジョブを全部やり直すには、ID のかわりに all を渡します。

bash
php artisan queue:retry all

失敗したジョブを消したいときは、queue:forget を使います。

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

補足

Horizon を使っているときは、queue:forget のかわりに horizon:forget で消します。

failed_jobs の表から、失敗したジョブを全部消すには、queue:flush を使います。

bash
php artisan queue:flush

queue:flush は、いつ失敗したかに関係なく、全部の記録を消します。--hours をつけると、指定した時間より前に失敗したものだけを消します。

bash
php artisan queue:flush --hours=48

消えたモデルを無視する#

ジョブにモデルを渡すと、列に入れるときに変換され、処理するときにデータベースから取り直されます。ところが、ジョブが順番を待っているあいだにモデルが消されると、ModelNotFoundException でジョブが失敗することがあります。

そのジョブを、例外を出さずに静かに捨てたいときは、PHP の属性 DeleteWhenMissingModels を付けます。

php
<?php

namespace App\Jobs;

use Illuminate\Queue\Attributes\DeleteWhenMissingModels;

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

失敗したジョブの記録を消す#

queue:prune-failed で、failed_jobs の古い記録を消せます。

bash
php artisan queue:prune-failed

既定では、24時間より古い記録が消えます。--hours をつけると、残す時間を決められます。次の例は、48時間より前に入った記録を消します。

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

失敗したジョブを DynamoDB に置く#

失敗したジョブの記録は、データベースの表のかわりに、DynamoDB にも置けます。ただし、記録を入れる DynamoDB の表は、自分で作る必要があります。表の名前は、ふつう failed_jobs です。アプリの queue 設定の queue.failed.table の値にそろえます。

failed_jobs の表には、文字列の主パーティションキー application と、文字列の主ソートキー uuid が要ります。application には、アプリの app 設定の name の値が入ります。アプリ名がキーに入るので、複数の Laravel アプリで同じ表を使えます。

Laravel から Amazon DynamoDB を使うために、AWS SDK も入れます。

bash
composer require aws/aws-sdk-php

次に、queue.failed.driver を dynamodb にします。失敗したジョブの設定に、key・secret・region も書きます。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 にします。ふつうは、環境変数 QUEUE_FAILED_DRIVER で決めます。

ini
QUEUE_FAILED_DRIVER=null

失敗したときのイベント#

ジョブが失敗したときに動くイベントリスナーを登録したいときは、Queue ファサードの failing を使います。たとえば、Laravel に最初から入っている AppServiceProvider の boot メソッドに、クロージャをつなげます。

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

列からジョブを消す#

補足

Horizon を使っているときは、queue:clear のかわりに horizon:clear で消します。

既定の接続の、既定のキューのジョブを全部消すには、queue:clear を使います。

bash
php artisan queue:clear

接続の引数と queue オプションを渡せば、特定の接続とキューのジョブを消せます。

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

注意

ジョブを消せるのは、SQS・Redis・データベースのキューだけです。また、SQS でメッセージを消すには最大60秒かかります。そのため、キューを空にしてから60秒以内に SQS に送ったジョブも、消えてしまうことがあります。

キューを見張る#

ジョブが急に増えると、キューがあふれて、ジョブの完了が遅くなることがあります。ジョブの数が決めた値を超えたら、Laravel に知らせさせられます。

まず、queue:monitor コマンドを、1分ごとに動かすようにします。見張りたいキューの名前と、ジョブの数の上限を渡します。

bash
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->connectionName,
                $event->queue,
                $event->size
            ));
    });
}

テスト#

ジョブを送るコードをテストするときは、ジョブそのものを実際には動かさないほうが、よいことが多いです。ジョブの中身は、送る側のコードとは別に、単独でテストできるからです。ジョブそのものをテストするには、ジョブのインスタンスを作って、テストの中で handle を直接呼びます。

Queue ファサードの fake を呼ぶと、ジョブが実際には列に入らなくなります。そのあと、アプリがジョブを列に入れようとしたかを、アサーション(「こうなっているはず」を確かめる命令)で調べられます。

Pest(PHP のテストの道具)で書く場合です。

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 exactly once...
    Queue::assertPushedOnce(ShipOrder::class);

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

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

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

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

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

PHPUnit(PHP のテストの道具)で書く場合です。

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 exactly once...
        Queue::assertPushedOnce(ShipOrder::class);

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

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

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

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

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

Queue のアサーションは、次のとおりです。

メソッド 説明
assertNothingPushed ジョブが1つも列に入っていないことを確かめる
assertPushedOn 決めたキューにジョブが入ったことを確かめる
assertPushed ジョブが列に入ったことを確かめる
assertPushedOnce ジョブがちょうど1回入ったことを確かめる
assertPushedTimes ジョブが決めた回数だけ入ったことを確かめる
assertNotPushed ジョブが入っていないことを確かめる
assertClosurePushed クロージャが列に入ったことを確かめる
assertClosureNotPushed クロージャが列に入っていないことを確かめる
assertCount 列に入ったジョブの総数を確かめる

assertPushed・assertNotPushed・assertClosurePushed・assertClosureNotPushed には、条件を書いたクロージャ(true か false を返す関数)を渡せます。条件に合うジョブが1つでも入っていれば、成功です。

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 に、にせものにするジョブのクラス名を渡します。

Pest で書く場合です。

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

PHPUnit で書く場合です。

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 ファサードのにせものの機能を使います。assertChained で、ジョブのチェーンが送られたかを確かめられます。第1引数に、つないだジョブの配列を渡します。

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

チェーンの変更をテストする#

チェーンの中のジョブが、チェーンの先頭や最後にジョブを足すなら、ジョブの assertHasChain で、残りのチェーンが期待どおりかを確かめられます。

php
$job = new ProcessPodcast;

$job->handle();

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

assertDoesntHaveChain は、残りのチェーンが空であることを確かめます。

php
$job->assertDoesntHaveChain();

チェーンの中のバッチをテストする#

チェーンの中にバッチがあるなら、チェーンのアサーションの中に 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 で、ジョブのバッチが送られたかを確かめられます。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 で、バッチが1つも送られていないことを確かめられます。

php
Bus::assertNothingBatched();

ジョブとバッチのやりとりをテストする#

1つのジョブが、そのバッチとどうやりとりしたかをテストしたいときもあります。たとえば、ジョブが、バッチの続きの処理を取り消したかを確かめたいときです。withFakeBatch で、にせもののバッチをジョブに割り当てます。withFakeBatch は、ジョブのインスタンスとにせもののバッチの組を返します。

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

$job->handle();

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

ジョブとキューのやりとりをテストする#

ジョブが、自分で列へ戻したか、自分を消したかを確かめたいときもあります。ジョブのインスタンスを作って、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();

テストに使えるメソッドをまとめます。

メソッド 説明
Bus::assertChained ジョブのチェーンが送られたことを確かめる
Bus::assertDispatchedWithoutChain チェーンなしでジョブが送られたことを確かめる
Bus::chainedBatch チェーンの中のバッチを調べる
assertHasChain ジョブの残りのチェーンが期待どおりか確かめる
assertDoesntHaveChain ジョブの残りのチェーンが空か確かめる
Bus::assertBatched バッチが送られたことを確かめる
hasJobs バッチに期待するジョブが入っているか確かめる
Bus::assertBatchCount 送られたバッチの数を確かめる
Bus::assertNothingBatched バッチが送られていないことを確かめる
withFakeBatch ジョブににせもののバッチを割り当てる
withFakeQueueInteractions ジョブとキューのやりとりをにせものにする
assertReleased ジョブが列へ戻されたことを確かめる(delay で待ち時間も)
assertDeleted ジョブが消されたことを確かめる
assertNotDeleted ジョブが消されていないことを確かめる
assertFailed ジョブが失敗にされたことを確かめる
assertFailedWith 決めた例外で失敗にされたことを確かめる
assertNotFailed ジョブが失敗にされていないことを確かめる

ジョブのイベント#

Queue ファサードの before と after を使うと、ジョブを処理する前と後に動くコールバックを決められます。記録を足したり、画面に出す数字を増やしたりするのに向いています。ふつうは、サービスプロバイダの 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 ファサードの looping では、ワーカーが列からジョブを取り出そうとする前に動くコールバックを決められます。たとえば、前に失敗したジョブが開きっぱなしにしたトランザクションを、取り消すのに使えます。

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

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

ワーカーが、列からジョブを取り出せなかったときは、Illuminate\Queue\Events\WorkerIdle というイベントも発生します。

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

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

関連するページ#

公式ドキュメント(英語)

2026年10月5日時点の内容をもとに、日本語でまとめています。

ページの一覧