RE:NODE

Базы данных12 мин чтения

Очереди задач на Valkey: BullMQ, Celery, RQ и Laravel

Фоновые задачи на Valkey с BullMQ, Celery, RQ или очередями Laravel: настройка, параметры, из-за которых теряются задачи, повторы, таймауты видимости и расчёт размера.

0 прочтений

Очередь задач выносит медленную работу из запроса: веб-процесс записывает небольшую запись вида «отправь это письмо» или «уменьши эту картинку», сразу же отвечает пользователю, а отдельный процесс-исполнитель забирает запись и выполняет работу. Чаще всего такую запись хранят в Valkey, потому что все популярные библиотеки очередей - BullMQ на Node, Celery и RQ на Python, очереди Laravel на PHP, Sidekiq на Ruby - построены на протоколе Redis, а Valkey говорит на нём без изменений.

Сложные вещи библиотеки делают хорошо: повторы, задержки, backoff, недоставленные задачи, параллелизм. Чего они не могут, так это защитить вас от Valkey, настроенного как кэш. Очередь на сервере, который вытесняет ключи при нехватке памяти, молча теряет задачи, а очередь на сервере без сохранения данных теряет все ожидающие задачи при перезапуске. В этой статье мы настроим каждую библиотеку, а затем разберём серверные настройки, сценарии отказов и расчёт размера, которые делают очередь надёжной.

Когда очередь - подходящий инструмент#

Очередь оправдывает себя, когда работа медленная, может упасть и быть повторена и не обязана закончиться до того, как пользователь увидит ответ. Обычные кандидаты:

  • Отправка писем, SMS и push-уведомлений.
  • Генерация миниатюр, PDF, выгрузок и отчётов.
  • Обращения к сторонним API, которые медленные или ограничивают частоту запросов.
  • Входящие webhooks, которые хочется сразу подтвердить, а обработать потом.
  • Работа по расписанию и повторяющаяся работа: ночные очистки, дайджесты, повторные попытки платежей.

Для чего она не подходит, так это для работы, которую пользователь ждёт. Если следующей странице нужен результат, очередь лишь добавит задержку и цикл опроса.

Прежде чем добавлять ради этого Valkey, проверьте, не справится ли уже ваша база данных. Таблица с SELECT ... FOR UPDATE SKIP LOCKED - корректная очередь в PostgreSQL и MySQL 8, и у неё есть свойство, которого нет ни у одной очереди на Valkey: постановка задачи и фиксация данных, которые её породили, происходят в одной транзакции. Этот вариант разобран в статье фоновые задачи на небольшом сервере. Valkey выигрывает, когда задач много, когда вам нужны возможности зрелой библиотеки, а не самописное решение, или когда Valkey у вас уже работает для сессий и кэша.

добавить задачублокирующий popмедленный вызовзапись результатаack или повторВеб-процессставит и отвечаетValkeyсписки и наборы задачБаза данныхнастоящие результатыСторонний сервиспочта, APIИсполнительзабирает и выполняет
Путь задачи от запроса до результата

Серверные настройки, от которых зависит очередь#

Три настройки на сервере Valkey решают, выживут ли задачи. Проверьте их до того, как писать код.

НастройкаЧто нужно очередиПочему
maxmemory-policynoevictionЛюбая другая политика может удалить ключи задач, когда память заполнена
appendonlyyesОжидающие задачи переживают перезапуск
appendfsynceverysecПри жёсткой остановке теряется не больше секунды постановок
maxmemoryЗадан, с запасомНакопившиеся задачи должны помещаться в память

Задачи теряет именно вытеснение. При allkeys-lru заполненный Valkey удаляет ключи, к которым дольше всего не обращались, а задача, дольше всех ждущая в наборе отложенных, - ровно такой ключ. BullMQ проверяет это при запуске и пишет в лог предупреждение, если политика не noeviction; остальные не проверяют вообще. С noeviction заполненный сервер отклоняет новые задачи с ошибкой, которую видит ваш код и на которую можно настроить оповещение. Это и есть нужное поведение: громкий отказ вместо тихого. Политики подробно разобраны в статье память и политики вытеснения Valkey.

bash
$ valkey-cli -h 203.0.113.20 -p 6380 --askpass INFO memory | grep -E 'maxmemory|used_memory_human'$ valkey-cli -h 203.0.113.20 -p 6380 --askpass INFO persistence | grep -E 'aof_enabled|rdb_last'

Если в одном Valkey живут кэш, которому положено вытесняться, и очередь, которой нельзя, честный ответ - два экземпляра. Маленький экземпляр для очереди обойдётся дешевле, чем отладка потерянных задач.

В RE:NODE тариф Valkey сохраняет данные на диск через AOF и снимки и защищён паролем с учётными данными, сгенерированными для каждого сервера. Он работает как отдельный сервер со своим хостом и портом, так что исполнители на тарифе для приложений и веб-процесс где-то ещё подключаются к одной и той же очереди.

BullMQ на Node.js#

BullMQ - актуальная библиотека очередей для Node, преемница Bull. Внутри она использует ioredis и хранит каждую очередь как набор ключей с префиксом bull:<queue>:.

bash
$ npm install bullmq
queue.js
import { Queue } from "bullmq";export const connection = {  host: process.env.VALKEY_HOST,  port: Number(process.env.VALKEY_PORT),  password: process.env.VALKEY_PASSWORD,};export const emails = new Queue("emails", {  connection,  defaultJobOptions: {    attempts: 5,    backoff: { type: "exponential", delay: 10_000 },    removeOnComplete: { age: 3600, count: 1000 },    removeOnFail: { age: 7 * 24 * 3600 },  },});// in the web processawait emails.add("welcome", { userId: 42 }, { jobId: `welcome:42` });
worker.js
import { Worker } from "bullmq";import { connection } from "./queue.js";const worker = new Worker("emails", async (job) => {  await sendWelcomeEmail(job.data.userId);}, { connection: { ...connection, maxRetriesPerRequest: null }, concurrency: 5 });worker.on("failed", (job, err) => console.error(job?.id, err.message));process.on("SIGTERM", async () => { await worker.close(); process.exit(0); });

Важные строки:

  • `maxRetriesPerRequest: null` в подключении исполнителя. Исполнители BullMQ используют блокирующие команды, которые ждут бесконечно; ioredis с лимитом повторов сдаётся на них и ломает исполнителя. BullMQ сам выставляет это значение, когда создаёт подключение, и предупреждает, если вы передаёте клиент с другой настройкой.
  • `removeOnComplete` и `removeOnFail`. По умолчанию BullMQ хранит каждую завершённую задачу. На загруженной очереди это самый частый способ заполнить Valkey: не ожидающей работой, а историей, которую никто не читает. Держите ограниченное окно.
  • `jobId`. Собственный id делает постановку идемпотентной: добавление задачи с уже существующим id ничего не делает. Так вы избегаете двойной отправки приветственного письма, когда запрос повторяется.
  • `worker.close()` по `SIGTERM`. Исполнитель перестаёт брать новые задачи и ждёт завершения текущих, так что deploy не прерывает задачу на середине. Почему этот обработчик сигнала важен на любой платформе, объясняет статья корректное завершение и health checks.

Повторяющиеся задачи - «каждую ночь в 03:00» - в актуальных версиях BullMQ встроены в виде планировщиков задач (job schedulers), а старый код использует параметр repeat у add. Используйте что-то одно, а не оба, и сверяйтесь с документацией для установленной версии.

Celery и RQ на Python#

Celery использует Valkey как брокер (где ждут сообщения) и, при желании, как хранилище результатов (где хранятся возвращаемые значения).

tasks.py
import osfrom celery import Celeryurl = os.environ["VALKEY_URL"]   # redis://:password@host:port/0app = Celery("shop", broker=url, backend=url.replace("/0", "/1"))app.conf.update(    task_acks_late=True,    worker_prefetch_multiplier=1,    task_reject_on_worker_lost=True,    result_expires=3600,    broker_transport_options={"visibility_timeout": 3600},)@app.task(bind=True, autoretry_for=(ConnectionError,), retry_backoff=True, max_retries=5)def send_welcome(self, user_id):    ...
bash
$ celery -A tasks worker --loglevel=INFO --concurrency=4$ celery -A tasks beat --loglevel=INFO

Настройка, в которой нужно разобраться, - visibility_timeout. С транспортом на протоколе Redis сообщение, взятое исполнителем, не удаляется, а скрывается до подтверждения. Если исполнитель не подтвердил его в пределах таймаута видимости - по умолчанию час, - сообщение доставляется снова другому исполнителю. Поэтому задача, которая законно выполняется дольше этого таймаута, выполнится дважды или много раз. Ставьте таймаут больше вашей самой длинной задачи, а с task_acks_late=True всё равно делайте задачи безопасными для повторного выполнения. result_expires важен по той же причине, что и removeOnComplete в BullMQ: результаты, которые вы никогда не читаете, копятся в виде ключей.

RQ - более простой выбор, когда вам не нужны маршрутизация, chords и canvas из Celery.

python
from redis import Redisfrom rq import Queue, Retryq = Queue("default", connection=Redis.from_url(os.environ["VALKEY_URL"]))job = q.enqueue(send_welcome, 42, retry=Retry(max=3, interval=[10, 60, 300]), job_timeout=600)
bash
$ rq worker --url "$VALKEY_URL" high default low

Имена очередей перечисляются в порядке приоритета: исполнитель опустошает high, прежде чем заглянуть в default. Упавшие задачи попадают в реестр упавших задач, который можно просмотреть и поставить задачи заново через rq info и панель RQ. Задачам по расписанию (enqueue_in, enqueue_at) нужен исполнитель, запущенный с --with-scheduler.

Очереди Laravel#

Системе очередей Laravel для работы на Valkey нужна только конфигурация.

.env
QUEUE_CONNECTION=redisREDIS_HOST=203.0.113.20REDIS_PORT=6380REDIS_PASSWORD=a-long-generated-password
bash
$ php artisan queue:work redis --queue=high,default --tries=3 --backoff=10 --max-time=3600

В config/queue.php у подключения redis есть значение retry_after, по умолчанию 90 секунд. Оно играет ту же роль, что и таймаут видимости в Celery: задача, зарезервированная дольше retry_after, возвращается в очередь. $timeout вашей самой длинной задачи должен быть на несколько секунд меньше retry_after, иначе длинные задачи выполняются дважды. Документация Laravel говорит об этом прямо, и всё равно это самая частая ошибка с очередями Laravel.

--max-time заставляет исполнителя чисто завершиться через час, чтобы тот, кто за ним присматривает, перезапустил его с нуля; это ограничивает утечки памяти в долгоживущих процессах PHP. После каждого deploy выполняйте php artisan queue:restart, потому что работающий исполнитель держит в памяти старый код. Laravel Horizon добавляет панель и автобалансировку поверх тех же очередей на протоколе Redis. Как присматривать за исполнителем, описано в статье deploy Laravel в продакшен.

Повторы, идемпотентность и недоставленные задачи#

Каждая очередь выше доставляет хотя бы один раз. Исполнитель может закончить задачу и умереть до подтверждения, и задача выполнится снова. Проектируйте с учётом этого, а не вопреки.

  1. Делайте задачи идемпотентными. «Отметить счёт 512 как оплаченный» можно безопасно повторить. «Прибавить 10 к балансу» - нельзя. Храните отметку об обработке - сильнее всего уникальное ограничение в базе данных на id события - и проверяйте её первым делом.
  2. Кладите в payload id, а не объекты. Задача, несущая целую запись пользователя, выполнится на устаревших данных, если прождёт час. Задача, несущая userId: 42, загрузит текущее состояние.
  3. Повторяйте с backoff, а потом остановитесь. Экспоненциальный backoff с пределом примерно в пять попыток справляется с временными сбоями. Задача, упавшая пять раз, не преуспеет на шестой; ей нужен человек.
  4. Держите упавшие задачи там, куда кто-то посмотрит. Набор failed в BullMQ, реестр упавших задач в RQ и таблица failed_jobs в Laravel - это очереди недоставленных задач (dead letters). Настраивайте оповещения на их размер, а не на отдельные отказы.
  5. Разделяйте медленные и быстрые очереди. Если ночная выгрузка и письма для сброса пароля делят одну очередь и двух исполнителей, выгрузка блокирует письма. Дайте срочной работе собственную очередь и исполнителей.

Расчёт размера Valkey для очереди#

Память очереди - это накопившиеся задачи плюс хранимая история. Типичная запись задачи - payload, параметры, временные метки и немного служебных данных библиотеки - занимает от нескольких сотен байт до нескольких килобайт. 100 000 накопившихся задач по 2 КБ - это около 200 МБ, и поэтому установившийся размер важен меньше, чем худший случай: что будет, если исполнители лежат час во время инцидента, а веб-процесс продолжает ставить задачи.

СитуацияЧто занимает памятьОриентир
Небольшое приложение, задачи выполняются за секундыМало; очередь почти пуста256 МБ с избытком
Ровный трафик, завершённые задачи хранятся часИстория завершённых задачЗадач в час, умноженных на размер
Большие payload (HTML, файлы в base64)Сами payloadХраните файл отдельно, ставьте в очередь ссылку
Исполнители остановлены из-за инцидентаВся накопившаяся очередьРассчитывайте на час постановок

Никогда не кладите в задачу файлы. Загрузите картинку на диск или в объектное хранилище и поставьте в очередь её путь. Измеряйте, а не прикидывайте: INFO memory для общего объёма, MEMORY USAGE на образце ключа задачи и собственные счётчики библиотеки (getJobCounts() в BullMQ, rq info, celery inspect).

Для очереди на Valkey CPU редко становится пределом; один маленький экземпляр выдерживает тысячи постановок в секунду. Исполнителям же нужны собственные CPU и память там, где они запущены.

За чем следить

Очередь отказывает медленно, а потом сразу целиком, поэтому следите за числами, которые сдвигаются первыми:

  • Размер очереди - ожидающие задачи в каждой очереди. Очередь, растущая днём и рассасывающаяся ночью, - это нормально; очередь, которая только растёт, означает, что исполнители не справляются.
  • Возраст самой старой ожидающей задачи. Полезнее, чем количество: 10 000 задач по миллисекунде - это ничто, десять задач, ждущих час, - это инцидент.
  • Упавшие задачи в час, с оповещением на рост, а не на отдельные отказы.
  • Память Valkey относительно лимита из INFO memory, с оповещением задолго до стены noeviction.
  • Подключённые клиенты из INFO clients. Парк исполнителей, теряющий подключения, упирается в лимит клиентов сервера.

У каждой библиотеки есть панель, показывающая большую часть этого, - Bull Board или Taskforce для BullMQ, Flower для Celery, rq-dashboard для RQ, Horizon для Laravel, - но панель, на которую никто не смотрит, - это не мониторинг. Экспортируйте два-три важных числа туда, откуда вам уже приходят оповещения; как их выбирать, рассказывает статья мониторинг, который что-то вам сообщает.

За исполнителями тоже нужно присматривать. Процесс-исполнитель, который завершился - необработанное исключение, остановка из-за нехватки памяти, deploy, - должен перезапускаться автоматически, иначе очередь тихо перестаёт рассасываться, а веб-процесс продолжает в неё добавлять. На VDS это юнит systemd с Restart=always; на хостинговой панели - собственное поведение сервера при перезапуске. В любом случае после каждого deploy проверяйте, что исполнители вернулись.

Устранение неполадок#

Задачи исчезают, не выполнившись и не упав. Вытеснение или перезапуск без сохранения данных. Проверьте evicted_keys в INFO stats и политику вытеснения.

Одна и та же задача выполняется дважды. Таймаут видимости или retry_after короче задачи, либо исполнитель умер до подтверждения. Увеличьте таймаут и сделайте задачу идемпотентной.

После сетевого сбоя исполнители перестают брать задачи. Подключение оборвалось, и клиент сдался. В BullMQ проверьте maxRetriesPerRequest: null; в остальных убедитесь, что за исполнителем присматривают и перезапускают его при выходе.

Память Valkey равномерно растёт, хотя очередь пуста. Завершённые задачи и результаты хранятся вечно. Задайте removeOnComplete, result_expires или аналог, а старые пусть истекут.

`OOM command not allowed when used memory > 'maxmemory'`. Политика noeviction делает свою работу. Очередь заполнена: исполнители лежат или слишком медленные, либо память заполняет история.

FAQ#

Может ли Valkey заменить RabbitMQ или Kafka?

Для фоновых задач одного приложения - да, и обслуживать придётся меньше. RabbitMQ лучше справляется со сложной маршрутизацией между множеством служб, а Kafka - с хранением больших журналов событий для повторного проигрывания. Если вы выбираете между ними для писем и миниатюр одного приложения, очередь на Valkey - соразмерный выбор.

Работают ли BullMQ, Celery и RQ с Valkey без изменений?

Они говорят на протоколе Redis и используют стандартные команды и Lua-скрипты, которые Valkey реализует. Вы подключаетесь тем же URL redis:// и теми же клиентскими библиотеками. Перед запуском проверьте свою версию на своём сервере, как при любом изменении инфраструктуры.

Должны ли очередь и кэш делить один Valkey?

Лучше не надо. Кэшу нужна политика с вытеснением, а очереди вытеснять нельзя никогда. На одном экземпляре придётся выбирать, и любой выбор будет неправильным для одной из них. Два маленьких экземпляра решают проблему.

Сколько исполнителей запускать?

Достаточно, чтобы в пиковые часы очередь оставалась почти пустой, и не больше, чем может загрузить работа. Для задач, упирающихся в ввод-вывод, таких как отправка писем, одного процесса с параллелизмом от 5 до 20 хватает надолго. Для задач, упирающихся в CPU, таких как обработка изображений, - по одному исполнителю на ядро.

Что будет с задачами по расписанию, если Valkey перезапустится?

При включённом сохранении они лежат на диске и загружаются вместе со всем остальным; отложенная задача, время которой прошло во время перезапуска, выполнится, когда исполнители переподключатся. Без сохранения все ожидающие и запланированные задачи пропадут.


Комментарии

Полностью анонимно: без аккаунта, без почты, без cookie. Мы храним имя, которое вы ввели, текст и время - больше ничего. Количество ссылок ограничено, разметка не отображается.

0/2000