Пользователь загружает резюме в PDF. Дальше нужно вытащить текст, разобрать его на структурированные поля, посчитать краткую сводку и сопоставить с открытыми вакансиями. На хорошем файле это 40 секунд, на скане с распознаванием — до трёх минут.
В HTTP-запросе такое не живёт. Первая версия жила и падала по таймауту nginx ровно на 60 секундах, причём файл к этому моменту был уже загружен, а профиль наполовину создан.
Три сообщения вместо одного
Очевидное решение — вынести всё в одну фоновую задачу. Я так и сделал сначала, и месяц спустя переписал на три.
final readonly class ProcessResumeMessage
{
public function __construct(public string $resumeId) {}
}
final readonly class GenerateSummaryMessage
{
public function __construct(public string $resumeId) {}
}
final readonly class MatchVacancyMessage
{
public function __construct(public string $resumeId, public ?string $vacancyId = null) {}
}Причина разделения не в чистоте архитектуры. Она в том, что у трёх этапов разные свойства отказа.
Разбор документа обращается к внешнему сервису извлечения текста и к распознаванию. Он долгий, дорогой и падает по сетевым причинам. Повторять его имеет смысл.
Сводка обращается к языковой модели. Тоже внешний вызов, но с другими лимитами и другой ценой ошибки: если сводки нет, профиль всё равно пригоден.
Сопоставление с вакансиями — чистый расчёт по своей базе. Оно быстрое, и его перезапускают отдельно каждый раз, когда появляется новая вакансия.
Слепив их в один обработчик, я получал одно из двух: либо повтор всей цепочки из-за сбоя на последнем шаге, либо ручное запоминание, докуда дошли. Второе — это самодельная машина состояний внутри обработчика, и она всегда пишется хуже, чем очередь.
Что передавать между шагами
Только идентификатор. Никакого текста документа, никаких разобранных полей.
Соблазн передать распознанный текст в следующее сообщение сильный: не надо второй раз читать из базы. В показанной конфигурации сообщение сериализуется и уезжает в Redis, а текст резюме — это от 5 до 200 килобайт. Умножаем на очередь из тысячи сообщений и получаем до двухсот мегабайт полезной нагрузки без учёта служебных данных.
Второй аргумент важнее. Сообщение с телом документа несёт снимок состояния на момент публикации. Если между шагами кто-то отредактировал профиль, следующий шаг работает с устаревшими данными и молча их перезаписывает. Идентификатор всегда указывает на текущее состояние.
#[AsMessageHandler]
final readonly class ProcessResumeMessageHandler
{
public function __invoke(ProcessResumeMessage $message): void
{
$resume = $this->resumes->find($message->resumeId)
?? throw new UnrecoverableMessageHandlingException('Резюме удалено');
$text = $this->parser->extract($resume->getFilePath());
$this->normalizer->fill($resume, $text);
$resume->setStage(ProcessingStage::Parsed);
$this->em->flush();
$this->bus->dispatch(new GenerateSummaryMessage($message->resumeId));
}
}Между flush() и dispatch() есть разрыв: если процесс умрёт в этой точке, стадия сохранится, а следующее сообщение не появится. Показанный фрагмент эту границу не закрывает. Надёжный вариант — outbox в одной транзакции с изменением стадии; более простой — периодический поиск записей, застрявших в промежуточном состоянии.
Отдельно про первую строку: если сущности нет, бросается UnrecoverableMessageHandlingException. Пользователь мог удалить резюме, пока оно стояло в очереди, и повторять обработку бессмысленно. Исключение отключает ретраи, но при настроенном failure transport сообщение всё равно окажется там; такие записи я отделяю от настоящих сбоев при разборе очереди ошибок.
Транспорты по весу задачи
Все три сообщения в одной очереди — плохая идея, которую я тоже успел попробовать. Одно резюме со сканом на 40 страниц занимает воркер на три минуты, и всё это время письма о регистрации стоят за ним.
Разведение по транспортам:
framework:
messenger:
transports:
heavy:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options: { queue_name: heavy }
retry_strategy: { max_retries: 5, delay: 5000, multiplier: 3 }
default:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options: { queue_name: default }
notifications:
dsn: '%env(MESSENGER_TRANSPORT_DSN)%'
options: { queue_name: notifications }
retry_strategy: { max_retries: 3, delay: 1000 }
routing:
'App\Message\ProcessResumeMessage': heavy
'App\Message\GenerateSummaryMessage': heavy
'App\Message\MatchVacancyMessage': default
'App\Message\NotificationMessage': notificationsГрупп воркеров три, а процессов четыре: два на тяжёлую очередь, один на обычную, один на уведомления. Тяжёлых больше, потому что они в основном ждут ответа от внешнего API, а не считают.
Ключевой момент, который я понял не сразу: воркер уведомлений должен быть отдельным процессом, а не просто отдельной очередью. Один воркер, подписанный на три транспорта, всё равно обрабатывает сообщения последовательно.
Внешний API внутри обработчика
Три правила, добытые опытом.
Нужны два таймаута. Symfony HttpClient берёт таймаут простоя из PHP-настройки default_socket_timeout, а общая длительность запроса по умолчанию не ограничена. Поэтому я задаю оба значения явно: поведение не зависит от окружения, а медленный ответ не занимает воркер бесконечно.
$response = $this->http->request('POST', $this->endpoint, [
'timeout' => 30,
'max_duration' => 120,
'body' => $payload,
]);timeout — простой без данных, max_duration — весь запрос целиком. Второй параметр забывают чаще, а именно он спасает от бесконечно медленного ответа.
Идемпотентность на стороне вызывающего. Повтор сообщения означает повторный вызов API. Если предыдущая попытка успела создать запись, а упала после этого, второй проход создаст дубликат. Одной проверки стадии мало: два воркера могут прочитать старое значение одновременно. Перед внешним вызовом обработчик атомарно переводит запись из ожидаемой стадии в Processing через условный UPDATE ... WHERE stage = :expected; ноль изменённых строк означает, что этап уже забрал другой процесс или он завершён.
Различать временные и постоянные ошибки. Статус ответа проверяется явно: сам вызов request() не обязан бросать исключение на HTTP-ошибке, а ClientExceptionInterface не покрывает 503. Для 429 и 503 я бросаю обычное исключение, чтобы сработал настроенный лимит повторов; RecoverableMessageHandlingException здесь не подходит, потому что принудительный повтор может обойти max_retries. Только известные ответы 400 и 422 означают неподдерживаемый документ.
$status = $response->getStatusCode();
if (in_array($status, [429, 503], true)) {
throw new \RuntimeException("Временная ошибка внешнего API: {$status}");
}
if (in_array($status, [400, 422], true)) {
$resume->fail(FailureReason::UnsupportedDocument);
$this->em->flush();
throw new UnrecoverableMessageHandlingException(
"Сервис отклонил документ: {$status}"
);
}Сетевые исключения обрабатываются отдельно как временные. Остальные неожиданные 4xx не маскируются под плохой файл: они уходят в журнал с телом ответа без персональных данных.
Логирование, которое отвечает на вопрос «почему висит»
Стандартный лог Messenger говорит, что сообщение получено и обработано. Этого мало, когда пользователь пишет «моё резюме обрабатывается уже двадцать минут».
Я завёл отдельный журнал этапов — не в файл, а в таблицу:
final class ResumeProcessingLogger
{
public function stage(string $resumeId, ProcessingStage $stage, array $context = []): void
{
$this->auditConnection->insert('resume_processing_log', [
'resume_id' => $resumeId,
'stage' => $stage->value,
'context' => json_encode($context, JSON_THROW_ON_ERROR),
'created_at' => (new \DateTimeImmutable())->format('Y-m-d H:i:s'),
]);
}
}Чтобы журнал действительно переживал откат доменных изменений, auditConnection должен быть отдельным DBAL-соединением в autocommit, а не соединением текущего EntityManager. Простого вызова DBAL недостаточно: на том же соединении вставка участвовала бы в активной транзакции. Для событий, которые обязаны пережить сбой базы целиком, нужен уже внешний журнал.
Теперь на вопрос про двадцать минут отвечает один запрос: видно, что документ разобран, сводка запрошена в 14:32, и с тех пор ничего. Дальше понятно, куда смотреть.
Прогресс для пользователя
Пользователю не нужен процент. Ему нужно понимать, что процесс идёт и когда примерно закончится.
Стадия хранится прямо в резюме энумом: загружено, разбирается, разобрано, сводка готова, сопоставлено, ошибка. Страница показывает её текстом и обновляется через опрос раз в три секунды.
Опрос вместо потока событий — сознательная экономия. Обработка идёт минуты, а не часы, пользователь смотрит на экран максимум пару минут. Городить хаб событий ради этого я не стал, хотя в другом проекте он есть и там оправдан.
Что оказалось важнее всего
Разделение на три сообщения. Всё остальное — техника, а это решение определило, насколько система понятна в проблемной ситуации.
Когда внешний сервис распознавания лежал полдня, у меня в очереди повторов накопилось около четырёхсот сообщений одного типа. Я видел ровно это: четыреста застрявших разборов, ноль проблем со сводками и сопоставлениями. После восстановления API достаточно было запустить consumer тяжёлого транспорта и дать ему разобрать накопившиеся повторы.
С монолитным обработчиком та же ситуация выглядела бы как четыреста провалившихся задач без понимания, на каком шаге и что с ними делать дальше.