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) для тяжёлых задач — иначе один воркер стоит, второй давится.
  • Задача может прийти дважды. Пишите обработчики идемпотентно.
Проверьте себя
1. Воркер получил задачу, начал её обрабатывать и упал через 3 секунды, не вызвав basic_ack. Что произойдёт с сообщением?
AОно потеряно — брокер удалил его при отправке
BБрокер вернёт его в очередь, и его подхватит другой воркер
CОно останется навсегда в состоянии unacked и заблокирует очередь
DБрокер отправит его в dead letter автоматически
2. Задачи в очереди тяжёлые: от 1 до 60 секунд. Какое значение prefetch_count разумно поставить воркеру?
Aprefetch_count=1 — следующую задачу берёт тот, кто освободился
Bprefetch_count=100 — чтобы воркер не простаивал в ожидании сети
Cprefetch не нужен, брокер и так раздаёт задачи справедливо
Dprefetch_count=0 — это и есть режим справедливой раздачи