Стандартная настройка повторов в Messenger выглядит разумно: три попытки, растущая задержка. Она перестаёт быть разумной в первый же день эксплуатации, потому что обращается одинаково с двумя совершенно разными событиями.
Внешний API вернул 503 — повторить через минуту правильно. Пользователь загрузил файл в неподдерживаемом формате — повторять бессмысленно, и три попытки просто откладывают на пятнадцать минут момент, когда он узнает об ошибке.
Три класса отказов
Я разложил всё, что падало за полгода, на три корзины.
Временный отказ. Внешний сервис недоступен, лимит запросов, разрыв соединения с базой, блокировка строки. Повторять нужно, и чем позже попытка, тем больше должна быть задержка.
Постоянный отказ данных. Файл повреждён, формат не поддерживается, обязательное поле пустое, сущность удалена. Повтор не изменит ничего. Такое сообщение должно уйти из очереди немедленно с внятной пометкой.
Отравленное сообщение. Обработчик падает с ошибкой в коде — опечатка в свойстве, несовпадение типов после рефакторинга. Это отдельный операционный класс: после исправления сообщение, возможно, обработается, но до релиза с исправлением каждый повтор только отравляет очередь.
Третий класс самый противный, потому что распознаётся только по симптому: одно и то же сообщение падает с одной и той же ошибкой на всех попытках, и ошибка — не сетевая.
Своя стратегия повторов
Messenger позволяет подменить стратегию целиком. Мне нужно было, чтобы задержка зависела от типа исключения.
final class CustomRetryStrategy implements RetryStrategyInterface
{
public function isRetryable(Envelope $message, ?\Throwable $throwable = null): bool
{
if ($throwable instanceof UnrecoverableExceptionInterface) {
return false;
}
$retries = RedeliveryStamp::getRetryCountFromEnvelope($message);
return $retries < $this->maxRetriesFor($throwable);
}
public function getWaitingTime(Envelope $message, ?\Throwable $throwable = null): int
{
$retries = RedeliveryStamp::getRetryCountFromEnvelope($message);
$delay = match (true) {
$throwable instanceof RateLimitExceededException => 60_000 * ($retries + 1),
$throwable instanceof TransportExceptionInterface => 5_000 * (3 ** $retries),
default => 10_000 * (2 ** $retries),
};
return random_int(
(int) ($delay * 0.8),
(int) ($delay * 1.2),
);
}
private function maxRetriesFor(?\Throwable $throwable): int
{
return $throwable instanceof RateLimitExceededException ? 8 : 4;
}
}Логика простая, но каждая строка выстрадана.
Превышение лимита внешнего API допускает восемь повторов с растущей базовой задержкой: лимиты обычно сбрасываются по окну, а агрессивный повтор через 5 секунд только усугубляет.
Для сетевых ошибок базовые интервалы растут втрое: 5, 15, 45, 135 секунд. Это четыре повтора после первоначальной доставки — достаточно, чтобы пережить короткий перезапуск соседнего сервиса.
Всё остальное получает до четырёх повторов с удвоением базовой задержки.
Случайный разброс в последних строках обязателен. Когда внешний сервис возвращается после падения, четыреста сообщений без jitter идут на повтор в одну секунду и роняют его снова. Я добавил разброс уже после такого повторного инцидента. В листинге взят диапазон ±20 процентов; это параметр примера, а не универсальный порог.
Dead letter в реальности
Failure transport в Messenger настраивается одной строкой, и на этом обычно останавливаются. Проблема начинается дальше: кто в неё смотрит.
У меня очередь проваленных задач читается двумя способами. Первый — команда, которую я запускаю руками, когда что-то расследую. Второй — обработчик, который срабатывает в момент окончательного провала:
#[AsEventListener(event: WorkerMessageFailedEvent::class)]
final readonly class FailedJobHandler
{
public function __invoke(WorkerMessageFailedEvent $event): void
{
if ($event->willRetry()) {
return;
}
$message = $event->getEnvelope()->getMessage();
$throwable = $event->getThrowable();
$this->logger->error('Задача провалена окончательно', [
'message' => $message::class,
'error' => $throwable->getMessage(),
]);
if ($message instanceof ProcessResumeMessage) {
$this->resumes->markFailed($message->resumeId, $throwable->getMessage());
}
$this->notifier->notifyOps($message, $throwable);
}
}Первая проверка обязательна. Без неё уведомление уходит на каждую промежуточную попытку, и через неделю на эти письма настраивается фильтр в почте — что равносильно отсутствию уведомлений.
Вторая важная часть — markFailed. Провалившаяся задача должна оставить след в предметной области, а не только в очереди. Пользователь заходит и видит «не удалось обработать: неподдерживаемый формат», а не бесконечное «обрабатывается».
Чистку я делаю раз в месяц руками. Автоматическое удаление проваленных сообщений через N дней выглядит аккуратно и убирает единственное свидетельство инцидента ровно тогда, когда до него дошли руки.
Четыре метрики
Из десятка возможных показателей очереди на дашборде нужны четыре.
Длина очереди. Очевидная и самая обманчивая. Ноль может означать «всё обработано» и «воркер умер и ничего не забирает» — во втором случае длина растёт, но если публикация тоже встала, она стоит на месте.
Возраст самого старого сообщения. Вот это главная метрика. Она отвечает на вопрос «сколько ждёт тот, кому не повезло больше всех». Длина в тысячу сообщений при возрасте в 20 секунд — нормальная нагрузка. Длина в пять при возрасте в час — авария.
Доля повторов. Отношение повторных доставок к общему числу. Скачок означает, что внешний сервис деградировал, ещё до того как задачи начнут проваливаться окончательно.
Число проваленных за сутки. Абсолютное значение, а не процент. Один провал в день — нормальный фон из битых файлов. Тридцать — что-то сломалось.
Реализация зависит от транспорта. Symfony Messenger хранит Redis-очередь в Streams, поэтому LLEN здесь не годится. Lag и pending надо брать из consumer group, например начать диагностику такими командами:
redis-cli XINFO GROUPS messages
redis-cli XPENDING messages symfonyДля Doctrine, AMQP и SQS источники будут другими. Я не прячу эту разницу за универсальным примером: общими остаются названия четырёх показателей, а не способ их получить.
Будит ночью, если настроено оповещение, только вторая метрика — возраст старейшего сообщения. Остальные три нужны, чтобы понять причину, когда уже проснулся.
Ручной повтор без дублирования
Кнопка «повторить» в админке — вещь, которая нужна каждый раз, когда чинишь последствия сбоя. И она же — способ создать сто дубликатов одним кликом.
Защита двухслойная. Первый слой: повтор возможен только для сущности в статусе ошибки, а статус меняется на «в обработке» до публикации. Между коммитом PostgreSQL и отправкой в Redis всё равно остаётся окно сбоя. Если нужна атомарная гарантия, в той же транзакции со статусом надо записывать outbox, а публиковать сообщение отдельным relay. Без outbox это две операции, а не одна транзакция.
Второй слой — проверка стадии в самом обработчике:
if ($resume->getStage()->isAtLeast(ProcessingStage::Parsed)) {
$this->logger->info('Разбор уже выполнен, пропускаю', ['resume' => $resume->getId()]);
return;
}Это делает обработчик идемпотентным по своей природе, а не по договорённости с интерфейсом. Дублирующее сообщение просто ничего не делает.
Массового повтора «перезапустить всё проваленное» у меня нет намеренно. Есть повтор по одному и есть консольная команда с явным фильтром по типу сообщения и по дате. Кнопка, которая одним нажатием отправляет в очередь четыреста задач, слишком легко нажимается.
Чему научил один инцидент
Внешний сервис распознавания начал отвечать медленно, но не падать: 40 секунд вместо трёх, при таймауте в 30. Каждый запрос заканчивался тайм-аутом, каждое сообщение уходило на повтор, повторы накапливались.
Длина очереди при этом выглядела прилично: воркеры активно работали. Первым делом я увидел не рост очереди, а рост доли повторов — и это единственная метрика, которая тогда сказала правду.
Вывод, который я оттуда унёс: очередь, которая работает, но ничего не завершает, выглядит здоровее пустой. Смотреть надо не на то, сколько сообщений забирают воркеры, а на то, сколько из них доходят до конца.