Блог AST-SoftPro
Асинхронная обработка: очереди, Celery, message brokers
Асинхронная обработка: очереди, Celery, message brokers
В бизнес-приложениях не все операции можно выполнить мгновенно. Отправка email, обработка изображений, генерация отчётов — всё это занимает время. Если выполнять такие операции в том же потоке, что и запрос пользователя, приложение будет зависать. Решение — асинхронная обработка через очереди задач.
Почему нужны очереди задач
Рассмотрим типичный сценарий: пользователь оформляет заказ. После сохранения заказа нужно:
-
Отправить email-подтверждение
-
Обновить счётчик товаров на складе
-
Отправить уведомление менеджеру в Telegram
-
Обновить аналитику
Если выполнять всё последовательно, запрос пользователя будет обрабатываться 3-5 секунд. С очередями — заказ сохраняется за 200 мс, а остальное выполняется в фоне.
Celery: стандарт для Python
Celery — распределённая система очередей задач для Python. Она работает с различными бэкендами (Redis, RabbitMQ) и поддерживает расписание, повторные попытки, группировку задач.
Установка и базовая настройка:
# celery_app.py
from celery import Celery
app = Celery(
"shop",
broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1",
)
app.conf.update(
task_serializer="json",
result_serializer="json",
accept_content=["json"],
timezone="Europe/Moscow",
enable_utc=True,
# Повторные попытки при сбое
task_acks_late=True,
worker_prefetch_multiplier=1,
)
Примеры задач:
# tasks.py
from celery_app import app
import requests
@app.task(bind=True, max_retries=3, default_retry_delay=60)
def send_order_email(self, order_id: int, email: str):
"""Отправка email с подтверждением заказа"""
try:
# Логика отправки email
requests.post("https://api.mailgun.net/v3/messages", json={
"to": email,
"subject": f"Заказ {order_id} оформлен",
"text": f"Ваш заказ {order_id} принят в обработку.",
})
except requests.RequestException as exc:
# Автоматический повтор через 60 секунд
raise self.retry(exc=exc)
@app.task
def process_image(image_url: str):
"""Обработка изображения: ресайз, компрессия"""
from PIL import Image
import io
response = requests.get(image_url)
img = Image.open(io.BytesIO(response.content))
img.thumbnail((800, 600))
# Сохранение обработанного изображения
return "processed"
@app.task
def generate_report(report_type: str, date_from: str, date_to: str):
"""Генерация отчёта (тяжёлая операция)"""
# Запрос к базе, агрегация, формирование PDF
return f"report_{report_type}_{date_from}_{date_to}.pdf"
@app.task
def notify_manager(order_id: int):
"""Уведомление менеджера в Telegram"""
requests.post(
"https://api.telegram.org/bot<TOKEN>/sendMessage",
json={
"chat_id": -1001234567890,
"text": f"Новый заказ #{order_id}",
},
)
Вызов задач из приложения:
from fastapi import FastAPI
from tasks import send_order_email, process_image, generate_report, notify_manager
app = FastAPI()
@app.post("/orders")
def create_order(order_data: dict):
# Быстрая операция — в основном потоке
order_id = save_order_to_db(order_data)
# Фоновые задачи — не блокируют ответ
send_order_email.delay(order_id, order_data["email"])
notify_manager.delay(order_id)
return {"id": order_id, "status": "created"}
@app.post("/reports")
def request_report(report_type: str, date_from: str, date_to: str):
# Запускаем генерацию в фоне
task = generate_report.delay(report_type, date_from, date_to)
return {"task_id": task.id, "status": "queued"}
@app.get("/reports/{task_id}")
def report_status(task_id: str):
from celery.result import AsyncResult
result = AsyncResult(task_id, app=celery_app.app)
return {
"status": result.status,
"result": result.result if result.ready() else None,
}
RabbitMQ vs Redis как брокер
Celery поддерживает разные брокеры. Выбор зависит от требований.
| Критерий | Redis | RabbitMQ |
|---|---|---|
| Надёжность доставки | Базовая (может потерять задачи при сбое) | Гарантия доставки (persistent queues) |
| Сложность настройки | Простая (уже есть для кэша) | Средняя (отдельный сервис) |
| Приоритеты задач | Нет | Да |
| Dead Letter Queue | Нет | Да |
| Делегирование по роутингу | Нет | Да (topic exchange) |
| Подходит для MVP | Да | Нет (избыточно) |
Когда нужен RabbitMQ:
-
Критические задачи (платежи, уведомления) — нельзя потерять
-
Нужны приоритеты (email важнее аналитики)
-
Нужен dead letter queue (неудачные задачи в отдельную очередь)
-
Сложный роутинг (разные воркеры для разных типов задач)
Настройка RabbitMQ с Celery:
app = Celery(
"shop",
broker="amqp://guest:guest@localhost:5672//",
backend="redis://localhost:6379/1",
)
app.conf.update(
task_routes={
"tasks.send_order_email": {"queue": "emails"},
"tasks.process_image": {"queue": "images"},
"tasks.generate_report": {"queue": "reports"},
"tasks.notify_manager": {"queue": "notifications"},
},
task_default_queue="default",
)
Запуск воркеров по очередям:
# Воркер для email (приоритетный)
celery -A celery_app worker -Q emails -c 4
# Воркер для изображений (ресурсоёмкий)
celery -A celery_app worker -Q images -c 2
# Воркер для отчётов (медленный)
celery -A celery_app worker -Q reports -c 1
Idempotency: защита от дублирования
В распределённых системах задача может выполниться дважды — например, если воркер упал после выполнения, но до отправки подтверждения. Idempotency гарантирует, что повторное выполнение не повредит данные.
Пример idempotent задачи:
@app.task(bind=True)
def process_payment(self, payment_id: int, amount: float):
"""Обработка платежа — idempotent"""
# Проверка: был ли платёж уже обработан
payment = Payment.query.get(payment_id)
if payment and payment.status == "completed":
return {"status": "already_completed", "id": payment_id}
try:
# Обработка платежа
result = payment_gateway.charge(payment_id, amount)
payment.status = "completed"
db_session.commit()
return {"status": "completed", "id": payment_id}
except Exception as exc:
raise self.retry(exc=exc, max_retries=3)
Паттерны idempotency:
-
Проверка статуса перед выполнением. Самый простой способ.
-
Idempotency key. Уникальный ключ в запросе, который гарантирует однократное выполнение.
-
Компенсационные транзакции. Если задача выполнилась частично — откат.
Расписание задач (Cron)
Celery Beat позволяет запускать задачи по расписанию.
from celery.schedules import crontab
app.conf.beat_schedule = {
# Ежедневная генерация отчёта
"daily-report": {
"task": "tasks.generate_report",
"schedule": crontab(hour=8, minute=0),
"args": ("sales", "today", "today"),
},
# Очистка устаревших сессий каждый час
"cleanup-sessions": {
"task": "tasks.cleanup_sessions",
"schedule": crontab(minute=0),
},
# Проверка просроченных заказов каждые 5 минут
"check-overdue-orders": {
"task": "tasks.check_overdue_orders",
"schedule": crontab(minute="*/5"),
},
}
Event-driven архитектура
Очереди задач — основа event-driven архитектуры. События генерируются при изменении состояния и обрабатываются подписчиками.
# Генерация события
@app.task
def order_created(order_id: int):
"""Событие: заказ создан"""
# Публикация события в очередь
app.send_task("tasks.send_order_email", args=[order_id])
app.send_task("tasks.notify_manager", args=[order_id])
app.send_task("tasks.update_analytics", args=[order_id])
# Подписчики обрабатывают событие независимо
Мониторинг очередей
Без мониторинга очереди задач — чёрный ящик. Что нужно отслеживать:
-
Длина очереди. Если очередь растёт — воркеры не справляются.
-
Время выполнения задач. Если время растёт — проблема в задаче или базе данных.
-
Количество ошибок. Если ошибки повторяются — задача неидемпотентна или есть проблема с ресурсами.
-
Uptime воркеров. Если воркер упал — задачи не выполняются.
# Flask-приложение для мониторинга
from fastapi import FastAPI
from celery import Celery
monitor_app = FastAPI()
@monitor_app.get("/queues")
def queue_stats():
inspect = celery_app.app.control.inspect()
return {
"active": inspect.active(),
"reserved": inspect.reserved(),
"scheduled": inspect.scheduled(),
}
@monitor_app.get("/workers")
def worker_stats():
inspect = celery_app.app.control.inspect()
return {
"ping": inspect.ping(),
"stats": inspect.stats(),
}
Частые ошибки
-
Синхронные задачи в очереди. Не кладите в очередь задачи, которые выполняются за 10 мс — накладные расходы на очередь больше, чем на выполнение.
-
Отсутствие idempotency. Платёж обработан дважды — это критическая ошибка. Всегда делайте задачи идемпотентными.
-
Одна очередь для всего. Разделяйте задачи по очередям: критические задачи (email, платежи) в отдельной очереди с приоритетом.
-
Отсутствие мониторинга. Без метрик вы не узнаете, что очередь забита, пока пользователи не начнут жаловаться.
-
Слишком много воркеров. 32 воркера на одном сервере — это не масштабирование, а борьба за ресурсы. Начинайте с 2-4 воркеров.
Заключение
Асинхронная обработка через очереди задач — обязательный компонент бизнес-приложений среднего и большого размера. Celery + Redis — хороший старт для большинства проектов. При росте требований переходите на RabbitMQ для надёжности и гибкости.
Главное правило: делайте задачи идемпотентными, разделяйте очереди по приоритетам и всегда мониторьте длину очереди.
Ключевые моменты:
-
Очереди задач (Celery) — для фоновой обработки длительных операций (email, отчёты, изображения).
-
Redis — простой брокер для MVP, RabbitMQ — для надёжности и приоритетов.
-
Idempotency обязательна для критических задач (платежи, уведомления).
-
Разделяйте задачи по очередям: критические и некритические отдельно.
-
Мониторьте длину очереди и время выполнения задач — без метрик очереди слепы.
Полезные ссылки
-
Celery — Distributed Task Queue — официальная документация
-
RabbitMQ Documentation — документация RabbitMQ
-
Redis — Open Source In-Memory Data Store — документация Redis
-
Idempotency in Distributed Systems — паттерны идемпотентности
-
Flower — Celery Monitoring — мониторинг Celery