HardТеория9 min

Очереди и задачи (Queues & Jobs)

Queue connections и drivers (sync, database, redis, sqs, beanstalkd), конфигурация очередей, создание и диспетчеризация задач, цепочки и пакеты задач

Очереди и задачи (Queues & Jobs)

Система очередей Laravel предоставляет унифицированный API для работы с различными бэкендами очередей: Amazon SQS, Redis, Beanstalkd, реляционные базы данных и синхронный драйвер для разработки. Очереди позволяют отложить выполнение ресурсоёмких задач, существенно улучшая время отклика приложения.

Конфигурация очередей

Конфигурация очередей находится в файле config/queue.php. Laravel 11 использует переменные окружения для настройки соединений.

// config/queue.php
return [
    'default' => env('QUEUE_CONNECTION', 'database'),

    'connections' => [
        'sync' => [
            'driver' => 'sync',
        ],

        'database' => [
            'driver' => 'database',
            'connection' => env('DB_QUEUE_CONNECTION'),
            'table' => env('DB_QUEUE_TABLE', 'jobs'),
            'queue' => env('DB_QUEUE', 'default'),
            'retry_after' => (int) env('DB_QUEUE_RETRY_AFTER', 90),
            'after_commit' => false,
        ],

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

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

        'beanstalkd' => [
            'driver' => 'beanstalkd',
            'host' => env('BEANSTALKD_HOST', 'localhost'),
            'queue' => env('BEANSTALKD_QUEUE', 'default'),
            'retry_after' => (int) env('BEANSTALKD_QUEUE_RETRY_AFTER', 90),
            'block_for' => 0,
            'after_commit' => false,
        ],
    ],

    'batching' => [
        'database' => env('DB_CONNECTION', 'sqlite'),
        'table' => 'job_batches',
    ],

    'failed' => [
        'driver' => env('QUEUE_FAILED_DRIVER', 'database-uuids'),
        'database' => env('DB_CONNECTION', 'sqlite'),
        'table' => 'failed_jobs',
    ],
];

Особенности драйверов

Sync - задачи выполняются синхронно в текущем процессе. Идеален для разработки и тестирования.

Database - требует миграцию для таблицы jobs. Подходит для небольших проектов без Redis.

Redis - наиболее популярный драйвер для продакшена. Требует расширение phpredis или пакет predis/predis.

SQS - Amazon Simple Queue Service. Не поддерживает delay больше 15 минут.

Beanstalkd - требует пакет pda/pheanstalk.

// Создание таблицы для database драйвера
// php artisan queue:table
// php artisan migrate

// Для пакетной обработки
// php artisan queue:batches-table
// php artisan migrate

Создание задач (Jobs)

Задачи создаются через Artisan-команду и размещаются в директории app/Jobs.

// php artisan make:job ProcessPodcast

<?php

declare(strict_types=1);

namespace App\Jobs;

use App\Models\Podcast;
use App\Services\AudioProcessor;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;
use Illuminate\Support\Facades\Log;

final class ProcessPodcast implements ShouldQueue
{
    use Dispatchable;
    use InteractsWithQueue;
    use Queueable;
    use SerializesModels;

    /**
     * The number of times the job may be attempted.
     */
    public int $tries = 3;

    /**
     * The maximum number of seconds the job can run.
     */
    public int $timeout = 120;

    /**
     * The number of seconds to wait before retrying.
     */
    public int $backoff = 10;

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

    public function __construct(
        public readonly Podcast $podcast,
    ) {}

    /**
     * Execute the job.
     */
    public function handle(AudioProcessor $processor): void
    {
        $processor->process($this->podcast);

        Log::info('Podcast processed successfully', [
            'podcast_id' => $this->podcast->id,
        ]);
    }

    /**
     * Handle a job failure.
     */
    public function failed(?\Throwable $exception): void
    {
        Log::error('Podcast processing failed', [
            'podcast_id' => $this->podcast->id,
            'error' => $exception?->getMessage(),
        ]);
    }
}

Уникальные задачи (Unique Jobs)

Уникальные задачи гарантируют, что в очереди одновременно будет только одна копия задачи с определённым ключом.

<?php

declare(strict_types=1);

namespace App\Jobs;

use App\Models\Product;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldBeUnique;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;

final class UpdateProductIndex implements ShouldQueue, ShouldBeUnique
{
    use Dispatchable;
    use InteractsWithQueue;
    use Queueable;
    use SerializesModels;

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

    public function __construct(
        public readonly Product $product,
    ) {}

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

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

    public function handle(): void
    {
        // Update search index for the product
    }
}

ShouldBeUniqueUntilProcessing

Важная разница: ShouldBeUnique держит блокировку до успешного завершения задачи, а ShouldBeUniqueUntilProcessing освобождает её, как только задача начинает выполняться.

<?php

declare(strict_types=1);

namespace App\Jobs;

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

final class SendReport implements ShouldQueue, ShouldBeUniqueUntilProcessing
{
    // Lock is released when the job starts processing,
    // allowing another identical job to be queued
}

Диспетчеризация задач (Dispatching Jobs)

Базовая диспетчеризация

<?php

declare(strict_types=1);

namespace App\Http\Controllers;

use App\Http\Requests\StorePodcastRequest;
use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Http\JsonResponse;

final class PodcastController extends Controller
{
    public function store(StorePodcastRequest $request): JsonResponse
    {
        $podcast = Podcast::create($request->validated());

        // Basic dispatch
        ProcessPodcast::dispatch($podcast);

        // Dispatch with delay
        ProcessPodcast::dispatch($podcast)
            ->delay(now()->addMinutes(10));

        // Dispatch to a specific queue
        ProcessPodcast::dispatch($podcast)
            ->onQueue('processing');

        // Dispatch to a specific connection
        ProcessPodcast::dispatch($podcast)
            ->onConnection('sqs');

        // Dispatch after database transaction commits
        ProcessPodcast::dispatch($podcast)
            ->afterCommit();

        // Dispatch synchronously (bypass queue)
        ProcessPodcast::dispatchSync($podcast);

        // Conditional dispatch
        ProcessPodcast::dispatchIf($podcast->isReady(), $podcast);
        ProcessPodcast::dispatchUnless($podcast->isProcessed(), $podcast);

        return response()->json($podcast, 201);
    }
}

Диспетчеризация через Facade

use Illuminate\Support\Facades\Queue;

// Push a job directly
Queue::push(new ProcessPodcast($podcast));

// Push to a specific queue
Queue::pushOn('processing', new ProcessPodcast($podcast));

// Push with a delay
Queue::later(now()->addMinutes(5), new ProcessPodcast($podcast));

Экспоненциальный backoff

<?php

declare(strict_types=1);

namespace App\Jobs;

use Illuminate\Contracts\Queue\ShouldQueue;

final class ProcessWebhook implements ShouldQueue
{
    /**
     * Calculate the number of seconds to wait before retrying.
     *
     * @return array<int, int>
     */
    public function backoff(): array
    {
        // Retry after 1, 5, 10 seconds
        return [1, 5, 10];
    }

    public function handle(): void
    {
        // Process webhook
    }

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

Цепочки задач (Job Chaining)

Цепочки позволяют выполнять задачи последовательно. Если одна задача в цепочке завершится неудачей, оставшиеся задачи не будут выполнены.

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

// Simple chain
Bus::chain([
    new ProcessPodcast($podcast),
    new OptimizePodcast($podcast),
    new ReleasePodcast($podcast),
    new NotifySubscribers($podcast),
])->dispatch();

// Chain with error callback
Bus::chain([
    new ProcessPodcast($podcast),
    new OptimizePodcast($podcast),
    new ReleasePodcast($podcast),
])->catch(function (\Throwable $e) {
    Log::error('Podcast chain failed', ['error' => $e->getMessage()]);
})->onQueue('podcasts')
  ->onConnection('redis')
  ->dispatch();

// Chain with closures
Bus::chain([
    new ProcessPodcast($podcast),
    function () use ($podcast) {
        $podcast->update(['status' => 'processed']);
    },
])->dispatch();

Пакеты задач (Job Batching)

Пакетная обработка позволяет запускать группу задач параллельно с обратными вызовами при завершении.

<?php

declare(strict_types=1);

namespace App\Jobs;

use Illuminate\Bus\Batchable;
use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;
use Illuminate\Queue\SerializesModels;

final class ImportCsvRow implements ShouldQueue
{
    use Batchable;
    use Dispatchable;
    use InteractsWithQueue;
    use Queueable;
    use SerializesModels;

    public function __construct(
        public readonly array $row,
    ) {}

    public function handle(): void
    {
        // Check if the batch has been cancelled
        if ($this->batch()?->cancelled()) {
            return;
        }

        // Process the CSV row
        \App\Models\Product::create($this->row);
    }
}

Создание пакета

use Illuminate\Bus\Batch;
use Illuminate\Support\Facades\Bus;
use App\Jobs\ImportCsvRow;

$rows = collect($csvData)->map(fn (array $row) => new ImportCsvRow($row));

$batch = Bus::batch($rows)
    ->then(function (Batch $batch) {
        // All jobs completed successfully
        Log::info('Import completed', ['batch_id' => $batch->id]);
    })
    ->catch(function (Batch $batch, \Throwable $e) {
        // First batch job failure detected
        Log::error('Import failed', ['error' => $e->getMessage()]);
    })
    ->finally(function (Batch $batch) {
        // The batch has finished executing (including failures)
        Notification::send(
            User::admins()->get(),
            new ImportCompleted($batch)
        );
    })
    ->allowFailures()
    ->onQueue('imports')
    ->name('CSV Import')
    ->dispatch();

// Access the batch ID
$batchId = $batch->id;

Работа с пакетом

use Illuminate\Support\Facades\Bus;

// Retrieve batch by ID
$batch = Bus::findBatch($batchId);

// Batch properties
$batch->id;                // Batch UUID
$batch->name;              // Batch name
$batch->totalJobs;         // Total number of jobs
$batch->pendingJobs;       // Number of jobs not yet processed
$batch->failedJobs;        // Number of failed jobs
$batch->processedJobs();   // Number of processed jobs
$batch->progress();        // Completion percentage (0-100)
$batch->finished();        // Whether all jobs are done
$batch->cancelled();       // Whether batch was cancelled

// Cancel the batch
$batch->cancel();

// Add more jobs to the batch
$batch->add([
    new ImportCsvRow($additionalRow),
]);

Пакеты с цепочками

Можно комбинировать пакеты и цепочки для сложных workflow.

use Illuminate\Support\Facades\Bus;

Bus::batch([
    // Each element is a chain
    [
        new DownloadImage($url1),
        new ProcessImage($url1),
        new UploadImage($url1),
    ],
    [
        new DownloadImage($url2),
        new ProcessImage($url2),
        new UploadImage($url2),
    ],
])->dispatch();

Запуск воркера очереди

# Basic worker
php artisan queue:work

# Specify connection and queue
php artisan queue:work redis --queue=high,default,low

# Process a single job
php artisan queue:work --once

# Process jobs for a specified duration
php artisan queue:work --stop-when-empty
php artisan queue:work --max-time=3600

# Maximum number of jobs before stopping
php artisan queue:work --max-jobs=1000

# Memory limit
php artisan queue:work --memory=256

# Sleep duration when no jobs available
php artisan queue:work --sleep=3

# Number of seconds to rest after exception
php artisan queue:work --backoff=3

# Number of times to attempt a job
php artisan queue:work --tries=3

Обработка проваленных задач

// View failed jobs
// php artisan queue:failed

// Retry a specific failed job
// php artisan queue:retry <uuid>

// Retry all failed jobs
// php artisan queue:retry all

// Delete a failed job
// php artisan queue:forget <uuid>

// Flush all failed jobs
// php artisan queue:flush

// Prune old failed jobs
// php artisan queue:prune-failed --hours=48

Программная обработка проваленных задач

<?php

declare(strict_types=1);

namespace App\Providers;

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

final class QueueServiceProvider extends ServiceProvider
{
    public function boot(): void
    {
        Queue::before(function (JobProcessing $event) {
            // $event->connectionName
            // $event->job
            // $event->job->payload()
        });

        Queue::after(function (JobProcessed $event) {
            // Job completed successfully
        });

        Queue::failing(function (JobFailed $event) {
            // $event->connectionName
            // $event->job
            // $event->exception
        });

        Queue::looping(function () {
            // Called each time the worker loop iterates
        });
    }
}

Laravel Horizon

Horizon предоставляет красивую панель мониторинга и конфигурацию для Redis-очередей.

// config/horizon.php (ключевые настройки)
'environments' => [
    'production' => [
        'supervisor-1' => [
            'maxProcesses' => 10,
            'balanceMaxShift' => 1,
            'balanceCooldown' => 3,
        ],
    ],

    'local' => [
        'supervisor-1' => [
            'maxProcesses' => 3,
        ],
    ],
],
// Notification for long wait times
use Laravel\Horizon\Horizon;

Horizon::routeMailNotificationsTo('[email protected]');
Horizon::routeSlackNotificationsTo('webhook-url', '#horizon');

Приоритизация очередей

<?php

declare(strict_types=1);

namespace App\Jobs;

use Illuminate\Bus\Queueable;
use Illuminate\Contracts\Queue\ShouldQueue;
use Illuminate\Foundation\Bus\Dispatchable;
use Illuminate\Queue\InteractsWithQueue;

final class SendInvoice implements ShouldQueue
{
    use Dispatchable;
    use InteractsWithQueue;
    use Queueable;

    public function __construct()
    {
        // Set the queue when constructing the job
        $this->onQueue('high');
    }

    public function handle(): void
    {
        // Process invoice
    }
}

// Worker processes queues in priority order:
// php artisan queue:work --queue=high,medium,default,low

Тестирование задач

<?php

declare(strict_types=1);

namespace Tests\Feature;

use App\Jobs\ProcessPodcast;
use App\Models\Podcast;
use Illuminate\Support\Facades\Bus;
use Illuminate\Support\Facades\Queue;
use Tests\TestCase;

final class PodcastTest extends TestCase
{
    public function test_podcast_is_dispatched(): void
    {
        Queue::fake();

        $podcast = Podcast::factory()->create();

        // Trigger the action that dispatches the job
        $this->postJson('/api/podcasts', [
            'title' => 'Test Podcast',
        ]);

        Queue::assertPushed(ProcessPodcast::class);

        Queue::assertPushed(ProcessPodcast::class, function ($job) use ($podcast) {
            return $job->podcast->id === $podcast->id;
        });

        Queue::assertPushedOn('processing', ProcessPodcast::class);

        Queue::assertNotPushed(AnotherJob::class);
    }

    public function test_job_chain_is_dispatched(): void
    {
        Bus::fake();

        // Trigger chain dispatch...

        Bus::assertChained([
            ProcessPodcast::class,
            OptimizePodcast::class,
            ReleasePodcast::class,
        ]);
    }

    public function test_batch_is_dispatched(): void
    {
        Bus::fake();

        // Trigger batch dispatch...

        Bus::assertBatched(function ($batch) {
            return $batch->jobs->count() === 10
                && $batch->jobs->every(fn ($job) => $job instanceof ImportCsvRow);
        });
    }
}

Проверь себя

Что делает опция after_commit в конфигурации очереди?

Что произойдёт, если задача в цепочке (Bus::chain) завершится неудачей?

Какой максимальный delay поддерживает драйвер SQS?

Какой метод Bus::batch() вызывается при ПЕРВОМ провале задачи в пакете?

Какой интерфейс нужно реализовать, чтобы задача была уникальной в очереди и блокировка снималась только после УСПЕШНОГО завершения?