Очереди и задачи (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);
});
}
}