MidПрактика5 min

Transports, Retry и Workers

Конфигурация транспортов: Doctrine, AMQP, Redis, in-memory, routing, retry strategy, failure transport, dead letter

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 для тестов

Проверь себя

Как направить одно сообщение в несколько транспортов одновременно?

Что произойдёт с сообщением, если handler выбросит исключение, реализующее `UnrecoverableExceptionInterface`?

В каком порядке worker обрабатывает сообщения из нескольких transports?

Какой transport лучше всего подходит для тестирования Messenger?

Сколько раз всего будет обработано сообщение при `max_retries: 3`?