Работа с очередями
Очереди нужны, когда работу не стоит выполнять в том же HTTP-запросе или artisan-команде, в которых она возникла. Типичные случаи — долгий пересчёт, пакетная миграция данных из мастер-сервисов, тяжёлое обращение к внешнему API, индексация пачки сущностей. Вы кладёте задачу в очередь и сразу возвращаете управление вызывающему коду; отдельный воркер позже подхватывает задачу и выполняет её в своём процессе.
В Goravel-сервисах Ensi очереди строятся на штатном Queue-фасаде фреймворка, а обвязка запуска воркеров (встроенный runner рядом с web и команда queue:work) лежит в app/adapters/queue. Сами джобы принадлежат модулям: app/modules/{module}/jobs. Список всех известных джоб собирается в bootstrap/jobs.go, настройки соединений и слушаемых очередей — в config/queue.go.
Очередь — это не Kafka. Kafka обычно несёт события интеграции между сервисами («случилось изменение в мастер-системе»). Очередь — внутренний механизм отложенной работы внутри сервиса: «сделай эту тяжёлую операцию позже, в воркере».
Как выглядит джоба со стороны фичи
Джоба — обычная структура с двумя обязательными методами. Signature() возвращает уникальное строковое имя задачи. Handle(args ...any) error содержит логику. Имя нужно воркеру, чтобы понять, какую структуру создавать и запускать: именно по signature задача лежит в Redis или в таблице очереди. Его лучше выбирать стабильным и говорящим, в духе catalog:reindex_offers или common:entities:migrate, и после появления в проде без крайней нужды не переименовывать — иначе уже отправленные, но ещё не выполненные задачи перестанут находиться.
Минимальный шаблон такой:
type ReindexOffersJob struct{}
func (j *ReindexOffersJob) Signature() string {
return "catalog:reindex_offers"
}
func (j *ReindexOffersJob) Handle(args ...any) error {
// аргументы приходят в том же порядке, в котором их передали при Dispatch
return nil
}
В реальных модулях джоба часто не держит всю бизнес-логику внутри Handle, а делегирует её action’у — так же, как HTTP-хэндлер не ходит в ORM напрямую. Например, джоба миграции сущностей создаёт (или получает) action и вызывает Execute. В Handle остаётся только обвязка: разобрать аргументы, вызвать действие, вернуть ошибку наружу.
Джобу обязательно зарегистрировать. Иначе воркер при разборе задачи из бэкенда не узнает signature и не сможет её выполнить. Удобный приём: в пакете jobs модуля собрать функцию All() []queue.Job, а в bootstrap/jobs.go склеить списки всех модулей в общее Jobs(). Новые джобы добавляют и в All() своего модуля, и тем самым — в общий перечень при старте приложения.
В шаблоне сервиса обычно оставляют учебные джобы в app/modules/examples/jobs (короткий и длинный sleep) — ими удобно проверить, что воркер жив, ещё до написания своей фичи.
Как отправить задачу
Из любого места, где уже поднято приложение — HTTP-хэндлер, artisan-команда, другая джоба, Kafka-handler — задачу отправляют через фасад:
err := facades.Queue().
Job(&jobs.ReindexOffersJob{}, []queue.Arg{
{Type: "int", Value: offerID},
{Type: "string", Value: "full"},
}).
OnQueue("default").
Dispatch()
Job принимает экземпляр джобы и слайс аргументов. OnQueue выбирает именованную очередь — логическую «ленту», из которой воркер будет забирать задачи. В типовом конфиге воркеры по умолчанию слушают очередь default. Если тяжёлые задачи не должны тормозить короткие, заводят отдельное имя (например heavy), кладут туда соответствующие джобы через OnQueue("heavy") и следят, чтобы воркер эту ленту реально слушал: либо добавляют имя в queue.listen, либо передают его аргументом в queue:work.
Аргументы описывают слайсом queue.Arg: у каждого элемента есть тип (string, int и другие значения, которые ожидает фреймворк) и значение. По типу значение сериализуется в бэкенд очереди и потом восстанавливается при вызове Handle. Порядок элементов при отправке совпадает с порядком в Handle(args ...any), поэтому в Handle обычно делают аккуратный type assert и не полагаются на «магический» разбор.
Иногда задачу нужно выполнить не сразу. Перед Dispatch() можно вызвать .Delay(time.Now().Add(5 * time.Minute)) — воркер возьмёт её только после указанного момента (для соединений, которые поддерживают отложенные задачи). Если в этом же процессе нужно дождаться результата и не класть ничего во внешнюю очередь, используют DispatchSync(): Handle выполнится синхронно в текущем потоке. Это полезно в тестах или в редких командах, где асинхронность мешает, но код отправки хочется оставить тем же.
Цепочку «сначала A, потом B» можно собрать через facades.Queue().Chain(...). Для большинства фич достаточно одной джобы с понятными аргументами.
Соединения: sync, redis, database
То, куда физически уходит задача после Dispatch(), задаёт QUEUE_CONNECTION (ключ queue.default в конфиге).
Если соединение sync (так часто в локальном .env.example), отдельный воркер не нужен: Dispatch() сразу вызывает Handle в том же процессе. Это удобно для отладки логики джобы «здесь и сейчас», но поведение не как в бою. Нет настоящего отложенного выполнения через бэкенд, нет отдельного воркера, падение джобы выглядит как ошибка вызывающего кода. Перед сдачей фичи, которая зависит от асинхронности, имеет смысл хотя бы раз прогнать сценарий на redis или database.
Для redis или database задача действительно уходит «в сторону» и должна быть забрана воркером. В конфиге для соединений задают драйвер, имя очереди по умолчанию и concurrent — сколько задач из одной очереди можно обрабатывать параллельно на одном воркере, если не переопределили это флагом --threads.
Неуспешные задачи попадают в хранилище failed jobs (в типовом конфиге — таблица failed_jobs в той же БД, что и сервис). Оттуда их можно разобрать и при необходимости переотправить штатными командами очереди Goravel (queue:failed, queue:retry и связанные).
Кто выполняет джобы
Dispatch только ставит задачу в бэкенд. Выполнить её должен воркер.
Первый способ — вместе с web-приложением. Если QUEUE_WORKER_ENABLED=true и текущее соединение не sync, при старте HTTP-процесса поднимается встроенный runner (app:queue) на очередях из queue.listen. Локально это удобно: подняли сервис — и короткие джобы уже кто-то обрабатывает. На stage/prod встроенный воркер обычно выключают и гоняют queue:work отдельным процессом или деплоем, чтобы масштаб воркеров не был жёстко привязан к числу HTTP-подов.
Второй способ — команда queue:work. Её запускают рядом с приложением или вместо встроенного runner:
./artisan queue:work
./artisan queue:work default
./artisan queue:work default heavy --threads=4
Без аргументов слушаются очереди из queue.listen. Имена после команды — конкретные очереди. Флаг --threads задаёт, сколько задач из одной очереди можно обрабатывать параллельно; если передать 0 или не указать флаг, берётся значение concurrent из конфига соединения (часто это 1).
В stdout для обрабатываемых задач выводятся статусы: RUNNING, затем DONE или FAIL и длительность.
Остановка — Ctrl+C локально или SIGTERM в Kubernetes. Текущие задачи стараются завершиться штатно: как и у Kafka-консьюмера, имеет смысл выставлять terminationGracePeriodSeconds с запасом на худшее время одного Handle, иначе при выкладке долгую джобу могут оборвать, и она уйдёт в failed / будет перезапущена по политике ретраев.
Ошибки и повторы
Если Handle вернул ошибку, джоба считается неуспешной. В прогресс-логе будет FAIL, запись появится среди failed jobs согласно queue.failed. Повторные попытки по умолчанию ограничены настройками фреймворка и соединения. Если политику ретраев нужно решать самим, у джобы реализуют ShouldRetry(err, attempt): метод говорит, стоит ли повторять выполнение и с какой паузой.
Как и для Kafka-хэндлеров, Handle лучше сразу писать идемпотентным. Воркер может взять задачу ещё раз после сбоя, рестарта пода или ручного queue:retry. Повтор не должен плодить дубли в БД, повторно слать одно и то же внешнее «разрушающее» действие без защиты или оставлять данные в полуприменённом состоянии без возможности дожать их следующим запуском.
Аргументы джобы тоже стоит проектировать с учётом повторов: передавайте устойчивые идентификаторы сущностей, а не «случайный снимок», который через пять минут уже не имеет смысла, если только это не осознанная часть контракта задачи.