Transport
Transport -- механизм доставки сообщений. Определяет, КАК сообщение передаётся от отправителя к обработчику.
Конфигурация транспортов
# config/packages/messenger.yaml
framework:
messenger:
transports:
# AMQP transport (RabbitMQ)
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
exchange:
name: messages
type: direct
queues:
messages:
binding_keys: [notification]
# High priority transport
high_priority:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options:
queues:
high_priority: ~
# Doctrine transport
doctrine_transport:
dsn: 'doctrine://default'
options:
queue_name: default
# Redis transport
redis_transport:
dsn: 'redis://localhost:6379/messages'
# Sync transport (for testing)
sync:
dsn: 'sync://'
# In-memory transport (for testing)
in_memory:
dsn: 'in-memory://'
routing:
'App\Message\SendNotification': async
'App\Message\ProcessOrder': high_priority
'App\Message\GenerateReport': async
'App\Message\CriticalAlert': sync
DSN-форматы
| Транспорт | DSN | Описание |
|---|---|---|
| AMQP | amqp://guest:guest@localhost:5672/%2f/messages |
RabbitMQ |
| Doctrine | doctrine://default?queue_name=default |
Таблица в БД |
| Redis | redis://localhost:6379/messages |
Redis Streams |
| Sync | sync:// |
Синхронная обработка |
| In-Memory | in-memory:// |
Только для тестов |
Подвох экзамена: Doctrine transport хранит сообщения в таблице
messenger_messages(по умолчанию). Он не требует отдельного брокера, но менее производителен, чем AMQP или Redis. Идеален для небольших проектов.
AMQP Transport (RabbitMQ)
framework:
messenger:
transports:
async:
dsn: 'amqp://guest:guest@localhost:5672/%2f'
options:
exchange:
name: app_messages
type: direct # direct, fanout, topic, headers
queues:
notifications:
binding_keys: ['notification']
orders:
binding_keys: ['order']
Exchange Types
| Тип | Описание |
|---|---|
direct |
Сообщение идёт в queue по routing key |
fanout |
Сообщение идёт во все привязанные queues |
topic |
Сообщение фильтруется по паттерну routing key |
headers |
Фильтрация по заголовкам сообщения |
Маршрутизация сообщений
Один transport
framework:
messenger:
routing:
'App\Message\SendNotification': async
Несколько transports
framework:
messenger:
routing:
# Message goes to BOTH transports
'App\Message\OrderPlaced': [async, audit]
Wildcard routing
framework:
messenger:
routing:
# All messages in this namespace go to async
'App\Message\*': async
# Override for specific messages
'App\Message\CriticalAlert': sync
Routing с Envelope stamps
<?php
declare(strict_types=1);
namespace App\Service;
use Symfony\Component\Messenger\Envelope;
use Symfony\Component\Messenger\MessageBusInterface;
use Symfony\Component\Messenger\Stamp\TransportNamesStamp;
final class SmartDispatcher
{
public function __construct(
private readonly MessageBusInterface $bus,
) {
}
public function dispatch(object $message, bool $highPriority = false): void
{
$stamps = [];
if ($highPriority) {
// Override configured transport
$stamps[] = new TransportNamesStamp(['high_priority']);
}
$this->bus->dispatch(new Envelope($message, $stamps));
}
}
Retry Strategy (Стратегия повторов)
Когда handler выбрасывает исключение, Messenger может автоматически повторить обработку.
framework:
messenger:
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
retry_strategy:
max_retries: 3
delay: 1000 # Initial delay: 1 second
multiplier: 2 # Exponential backoff
max_delay: 60000 # Max delay: 60 seconds
high_priority:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
retry_strategy:
max_retries: 5
delay: 500
multiplier: 3
max_delay: 120000
Exponential backoff
Attempt 1: delay = 1000ms (1s)
Attempt 2: delay = 1000 * 2 = 2000ms (2s)
Attempt 3: delay = 2000 * 2 = 4000ms (4s)
...
Attempt N: delay = min(calculated, max_delay)
Подвох экзамена:
max_retries: 3означает 3 повторных попытки ПОСЛЕ первого неудачного вызова. Итого сообщение будет обработано максимум 4 раза (1 оригинальный + 3 retry).
UnrecoverableExceptionInterface
<?php
declare(strict_types=1);
namespace App\Exception;
use Symfony\Component\Messenger\Exception\UnrecoverableExceptionInterface;
// Messages throwing this exception will NOT be retried
final class InvalidOrderException
extends \RuntimeException
implements UnrecoverableExceptionInterface
{
public static function notFound(int $orderId): self
{
return new self(sprintf('Order #%d not found', $orderId));
}
}
| Интерфейс | Поведение |
|---|---|
UnrecoverableExceptionInterface |
Никогда не retry, сразу в failure transport |
RecoverableExceptionInterface |
Всегда retry (даже если max_retries = 0) |
| Обычное исключение | Retry согласно стратегии |
Failure Transport (Dead Letter Queue)
framework:
messenger:
# Global failure transport
failure_transport: failed
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
retry_strategy:
max_retries: 3
delay: 1000
multiplier: 2
# Failure transport stores permanently failed messages
failed:
dsn: 'doctrine://default?queue_name=failed'
Жизненный цикл сообщения с ошибкой
1. Handler throws exception
2. Retry strategy: retries left?
YES -> Re-queue with delay -> Go to step 1
NO -> Send to failure transport
3. Message sits in failure transport
4. Admin reviews and decides:
- Retry (messenger:failed:retry)
- Remove (messenger:failed:remove)
Команды для работы с failed messages
# Show all failed messages
php bin/console messenger:failed:show
# Retry a specific failed message
php bin/console messenger:failed:retry 42
# Retry ALL failed messages
php bin/console messenger:failed:retry
# Remove a specific failed message
php bin/console messenger:failed:remove 42
Per-transport Failure Transport
framework:
messenger:
failure_transport: failed_default
transports:
async:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
failure_transport: failed_async # Specific failure transport
high_priority:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
failure_transport: failed_priority
failed_default:
dsn: 'doctrine://default?queue_name=failed_default'
failed_async:
dsn: 'doctrine://default?queue_name=failed_async'
failed_priority:
dsn: 'doctrine://default?queue_name=failed_priority'
Workers (Воркеры)
Worker -- процесс, который потребляет сообщения из transport и передаёт их handlers.
# Consume from one transport
php bin/console messenger:consume async
# Consume from multiple transports (with priority)
php bin/console messenger:consume high_priority async
# With time limit
php bin/console messenger:consume async --time-limit=3600
# With message limit
php bin/console messenger:consume async --limit=100
# With memory limit
php bin/console messenger:consume async --memory-limit=256M
Параметры Worker
| Параметр | Описание |
|---|---|
--limit=N |
Обработать N сообщений и остановиться |
--time-limit=N |
Работать N секунд и остановиться |
--memory-limit=N |
Остановиться при достижении лимита памяти |
--sleep=N |
Ожидание между проверками очереди (секунды) |
--queues=Q |
Потреблять только из указанных очередей |
Worker проверяет transports в указанном порядке. Если в high_priority есть сообщения, они обрабатываются первыми.
Worker в Production (Supervisor)
; /etc/supervisor/conf.d/messenger-worker.conf
[program:messenger-consume]
command=php /path/to/project/bin/console messenger:consume async high_priority --time-limit=3600
user=www-data
numprocs=2
startsecs=0
autostart=true
autorestart=true
startretries=10
process_name=%(program_name)s_%(process_num)02d
Тестирование Messenger
# config/packages/messenger.yaml (test environment)
when@test:
framework:
messenger:
transports:
async: 'in-memory://'
high_priority: 'in-memory://'
<?php
declare(strict_types=1);
namespace App\Tests\Service;
use App\Message\SendNotification;
use Symfony\Bundle\FrameworkBundle\Test\KernelTestCase;
use Zenstruck\Messenger\Test\InteractsWithMessenger;
final class OrderServiceTest extends KernelTestCase
{
use InteractsWithMessenger;
public function testOrderPlacedDispatchesNotification(): void
{
// ... place order ...
$this->transport('async')
->queue()
->assertContains(SendNotification::class, 1);
// Process messages
$this->transport('async')->process();
$this->transport('async')
->queue()
->assertEmpty();
}
}
Итоги
- Transport определяет механизм доставки: AMQP, Doctrine, Redis, sync, in-memory
- Routing связывает тип сообщения с transport; wildcard
*для namespace - Retry strategy:
max_retries,delay,multiplier,max_delay max_retries: 3= 4 попытки всего (1 + 3 retry)UnrecoverableExceptionInterface-- сразу в failure, без retry- Failure transport хранит окончательно провалившиеся сообщения
- Worker (
messenger:consume) потребляет сообщения; Supervisor для production in-memory://transport для тестов