Work Queue: очередь задач
Тяжёлую работу нельзя делать внутри HTTP-запроса — её надо отдать воркерам. Разбираем самый частый паттерн RabbitMQ: очередь задач.
Work Queue (очередь задач) — одна очередь, много одинаковых консьюмеров-воркеров. Каждое сообщение достаётся ровно одному из них, и брокер сам следит, чтобы задача не потерялась, если воркер упал.
Зачем это нужно
Пользователь загружает фотографию в профиль. Её надо привести к четырём размерам, прогнать через оптимизатор, залить в хранилище. На это уходит секунд восемь. Если делать это прямо в обработчике HTTP-запроса, пользователь восемь секунд смотрит на крутилку, а веб-сервер держит воркер gunicorn занятым — и на десяти параллельных загрузках приложение ложится.
Правильный ход: обработчик кладёт в очередь маленькое сообщение «сделай ресайз картинки 4821» и мгновенно отвечает 202 Accepted. Дальше картинку разбирает отдельный процесс. Пользователь свободен, веб-сервер свободен, а если очередь растёт — вы запускаете ещё три копии воркера, и она рассасывается. Это и есть главная выгода паттерна: масштабирование горизонтально, без переписывания кода.
Продюсер: кто ставит задачу
import json
import pika
params = pika.ConnectionParameters(
host="localhost",
credentials=pika.PlainCredentials("guest", "guest"),
heartbeat=60,
)
connection = pika.BlockingConnection(params)
channel = connection.channel()
# durable=True — описание очереди переживёт перезапуск брокера
channel.queue_declare(queue="image.resize", durable=True)
task = {
"image_id": 4821,
"source": "s3://uploads/4821/original.jpg",
"sizes": [128, 256, 512, 1024],
}
channel.basic_publish(
exchange="", # default exchange
routing_key="image.resize", # ...= имя очереди
body=json.dumps(task).encode(),
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent, # сообщение на диск
content_type="application/json",
),
)
print("задача 4821 поставлена в очередь")
connection.close()
Разберём по строкам то, что обычно пропускают.
exchange=""— это default exchange, служебный обменник, к которому каждая очередь привязана под своим именем. Поэтомуrouting_key="image.resize"означает буквально «положи в очередь image.resize». Никакой магии, просто удобное сокращение.durable=Trueсохраняет описание очереди, аdelivery_mode=Persistent— сами сообщения. Это два разных флага, и нужны оба: durable-очередь с транзиентными сообщениями после перезапуска брокера окажется пустой.- В теле —
image_idи ссылка на файл, а не сами байты картинки. Сообщение должно быть маленьким: очередь — это транспорт указаний, а не файловое хранилище.
Воркер: кто задачу делает
import json
import time
import pika
def resize(task):
# здесь реальная работа: скачать, порезать, залить
time.sleep(2)
if task["image_id"] % 13 == 0:
raise ValueError("битый JPEG")
def on_message(ch, method, properties, body):
task = json.loads(body)
print("взял", task["image_id"], flush=True)
try:
resize(task)
except Exception as exc:
print("не смог:", exc, flush=True)
# requeue=False — в dead letter, а не по кругу
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
return
ch.basic_ack(delivery_tag=method.delivery_tag)
print("готово", task["image_id"], flush=True)
connection = pika.BlockingConnection(
pika.ConnectionParameters(host="localhost", heartbeat=60)
)
channel = connection.channel()
channel.queue_declare(queue="image.resize", durable=True)
channel.basic_qos(prefetch_count=1) # не бери следующую, пока не сдал текущую
channel.basic_consume(
queue="image.resize",
on_message_callback=on_message,
auto_ack=False, # подтверждаем руками
)
print("жду задачи, Ctrl+C для выхода")
channel.start_consuming()
Запустите этот файл в двух-трёх терминалах — вот и всё масштабирование. Брокер сам раздаст сообщения между подключёнными консьюмерами. Ни продюсер, ни другие воркеры об этом не узнают: очередь остаётся той же, меняется только число ртов, которые из неё едят.
Как это работает
Подтверждения (ack) и почему задача не теряется
Когда брокер отдал сообщение воркеру, он его не удаляет — помечает как «в работе» (unacked) и ждёт. Удаление происходит только после basic_ack. Если воркер умер на середине — отвалилась сеть, кто-то сделал kill -9, кончилась память — TCP-соединение рвётся, брокер видит это и возвращает unacked-сообщение в очередь. Его подхватит другой воркер. Никакого таймаута задачи тут нет: RabbitMQ ждёт ровно столько, сколько живо соединение (по умолчанию, если не настроен consumer_timeout).
Отсюда важное следствие: задачи обязаны быть идемпотентными. Сообщение может быть доставлено дважды — например, воркер успел залить все четыре превью, но упал за миг до ack. Повторный ресайз безобиден, а вот повторное списание денег — нет.
prefetch_count: почему без него всё тормозит
Без basic_qos брокер вываливает воркеру сообщения по кругу, не глядя на то, занят тот или нет. В итоге медленные задачи скапливаются у одного воркера, а другой стоит без дела. С prefetch_count=1 следующая задача уходит только тому, кто освободился. Разница видна на арифметике — этот пример можно запустить прямо здесь:
tasks = [10, 1, 1, 1, 10, 1] # секунд на задачу
# слепая раздача по кругу: 1-я, 3-я, 5-я — первому воркеру и т.д.
rr = [0, 0]
for i, cost in enumerate(tasks):
rr[i % 2] += cost
# prefetch_count=1: следующую берёт тот, кто освободился
busy = [0, 0]
for cost in tasks:
w = busy.index(min(busy))
busy[w] += cost
print("round-robin заранее:", rr, "-> всё готово через", max(rr), "c")
print("prefetch_count=1: ", busy, "-> всё готово через", max(busy), "c")
Результат:
round-robin заранее: [21, 3] -> всё готово через 21 c
prefetch_count=1: [11, 13] -> всё готово через 13 c
Одни и те же шесть задач, одни и те же два воркера — и восемь секунд разницы просто из-за одной строчки настройки. На реальной нагрузке разрыв больше.
Сколько ставить prefetch
| Задачи | prefetch_count | Почему |
| Тяжёлые, секунды и минуты (ресайз, отчёты) | 1 | Справедливая раздача важнее экономии на round-trip |
| Лёгкие, миллисекунды (отправка пуша) | 50–200 | Иначе воркер простаивает в ожидании сети |
| Смешанные | Разные очереди | Не смешивайте минутные задачи с миллисекундными в одной очереди |
Что делать с упавшей задачей
Отказы бывают двух сортов, и путать их дорого. Временный отказ — S3 не ответил, база ушла в перезагрузку. Такую задачу имеет смысл повторить. Постоянный — файл битый, формат не поддерживается. Сколько ни повторяй, результат один.
Проблема в том, что basic_nack(requeue=True) возвращает сообщение в голову очереди немедленно. Воркер тут же берёт его снова, снова падает — и вы получаете цикл на тысячи итераций в секунду, который выжигает CPU и забивает логи. Поэтому мгновенный requeue почти всегда неправильный ответ.
Рабочая схема: у очереди настроен dead letter exchange (DLX), а отказ отправляется с requeue=False. Сообщение уходит в отдельную очередь image.resize.dead, где спокойно ждёт разбора — вы видите его в мониторинге, читаете причину, чините код и переливаете задачи обратно. Если нужны именно повторы, делают очередь ожидания с TTL: сообщение падает в неё на 30 секунд, по истечении TTL DLX возвращает его в рабочую очередь. Так получается retry с задержкой, а счётчик попыток кладут в заголовки сообщения — после третьей попытки задача едет в dead letter насовсем.
Частые ошибки
auto_ack=True«чтобы не возиться». Сообщение считается доставленным в момент отправки. Воркер упал — задача исчезла молча, и вы узнаете об этом от пользователя, у которого не появилась аватарка. Автоподтверждение допустимо только там, где потеря не страшна (метрики, логи).- Ack в начале обработки. Логически это тот же
auto_ack, только написанный вручную. Подтверждайте после того, как работа доделана и результат сохранён. - Долгая задача рвёт соединение.
BlockingConnectionв pika однопоточное: пока внутри вашего колбэка крутится десятиминутный ресайз, библиотека не отвечает на heartbeat, и брокер закрывает соединение с ошибкойConnectionResetError. Лечится либо периодическимconnection.process_data_events()внутри длинной работы, либо выносом работы в отдельный поток и подтверждением черезconnection.add_callback_threadsafe(...). requeue=Trueпри отказе. Битая картинка вернётся в очередь, снова упадёт, снова вернётся — получится бесконечный цикл, который сожжёт CPU всех воркеров. Отправляйте такие сообщения в dead letter (requeue=Falseплюс настроенный DLX) и разбирайте отдельно.- Картинка целиком в теле сообщения. Мегабайтные сообщения раздувают память брокера и делают его узким местом. Кладите ссылку на объектное хранилище.
- Нет ограничения на рост очереди. Если продюсеры быстрее консьюмеров, очередь растёт до тех пор, пока брокер не упрётся в память и не включит flow control. Следите за глубиной очереди в мониторинге — это ваш главный сигнал «пора добавить воркеров».
Итоги
- Work Queue = одна очередь + N одинаковых воркеров; сообщение достаётся ровно одному.
- Масштабирование — просто запуск ещё одного процесса-консьюмера, код не меняется.
durable=Trueдля очереди иdelivery_mode=Persistentдля сообщений — оба флага, иначе перезапуск брокера съест задачи.auto_ack=False+ ручнойbasic_ackпосле работы: упавший воркер вернёт задачу в очередь.basic_qos(prefetch_count=1)для тяжёлых задач — иначе один воркер стоит, второй давится.- Задача может прийти дважды. Пишите обработчики идемпотентно.