Подписывайтесь:

Блог AST-SoftPro

Асинхронная обработка: очереди, Celery, message brokers

05.07.2026 13 мин чтения
Асинхронная обработка: очереди, Celery, message brokers

Асинхронная обработка: очереди, Celery, message brokers

В бизнес-приложениях не все операции можно выполнить мгновенно. Отправка email, обработка изображений, генерация отчётов — всё это занимает время. Если выполнять такие операции в том же потоке, что и запрос пользователя, приложение будет зависать. Решение — асинхронная обработка через очереди задач.

Почему нужны очереди задач

Рассмотрим типичный сценарий: пользователь оформляет заказ. После сохранения заказа нужно:

  1. Отправить email-подтверждение

  2. Обновить счётчик товаров на складе

  3. Отправить уведомление менеджеру в Telegram

  4. Обновить аналитику

Если выполнять всё последовательно, запрос пользователя будет обрабатываться 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:

  1. Проверка статуса перед выполнением. Самый простой способ.

  2. Idempotency key. Уникальный ключ в запросе, который гарантирует однократное выполнение.

  3. Компенсационные транзакции. Если задача выполнилась частично — откат.

Расписание задач (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(),
    }

Частые ошибки

  1. Синхронные задачи в очереди. Не кладите в очередь задачи, которые выполняются за 10 мс — накладные расходы на очередь больше, чем на выполнение.

  2. Отсутствие idempotency. Платёж обработан дважды — это критическая ошибка. Всегда делайте задачи идемпотентными.

  3. Одна очередь для всего. Разделяйте задачи по очередям: критические задачи (email, платежи) в отдельной очереди с приоритетом.

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

  5. Слишком много воркеров. 32 воркера на одном сервере — это не масштабирование, а борьба за ресурсы. Начинайте с 2-4 воркеров.

Заключение

Асинхронная обработка через очереди задач — обязательный компонент бизнес-приложений среднего и большого размера. Celery + Redis — хороший старт для большинства проектов. При росте требований переходите на RabbitMQ для надёжности и гибкости.

Главное правило: делайте задачи идемпотентными, разделяйте очереди по приоритетам и всегда мониторьте длину очереди.

Ключевые моменты:

  1. Очереди задач (Celery) — для фоновой обработки длительных операций (email, отчёты, изображения).

  2. Redis — простой брокер для MVP, RabbitMQ — для надёжности и приоритетов.

  3. Idempotency обязательна для критических задач (платежи, уведомления).

  4. Разделяйте задачи по очередям: критические и некритические отдельно.

  5. Мониторьте длину очереди и время выполнения задач — без метрик очереди слепы.

Полезные ссылки

AI-Помощник