Очереди сообщений в Python: практическое руководство по Celery, RabbitMQ и Redis | AdminWiki

Очереди сообщений в Python: практическое руководство по Celery, RabbitMQ и Redis

18 августа 2026 7 мин. чтения

Введение в очереди сообщений и Celery

Очередь сообщений решает задачу асинхронной обработки данных. Вместо того чтобы заставлять пользователя ждать завершения длительной операции, приложение отправляет задачу в очередь и сразу отвечает. Брокер хранит задачу, а воркер забирает её и выполняет в фоне. Это снижает задержку ответа и повышает отказоустойчивость системы.

Celery - это распределённая очередь задач для Python. Библиотека берёт на себя управление задачами: отправку, маршрутизацию, повторные попытки и хранение результатов. Celery работает с брокерами сообщений RabbitMQ и Redis, которые отвечают за передачу задач от приложения к воркерам. RabbitMQ обеспечивает надёжную доставку и сложную маршрутизацию, Redis - высокую скорость и простоту настройки.

Типовые сценарии использования Celery: отправка email после регистрации пользователя, обработка изображений, генерация отчётов, синхронизация данных с внешними API. Любая операция, которая занимает больше пары секунд, кандидат на вынос в фоновую задачу. Это разгружает веб-сервер и делает интерфейс отзывчивым.

Базовые принципы работы очередей сообщений разобраны в практическом руководстве по оптимизации логирования в Python. Для выбора между RabbitMQ и Redis полезно изучить алгоритм выбора брокера сообщений.

Установка и настройка брокера сообщений

Первый шаг - установить брокер. От выбора зависит надёжность доставки, скорость и сложность эксплуатации.

Выбор брокера: RabbitMQ или Redis?

КритерийRabbitMQRedis
Надёжность доставкиВысокая, подтверждения publisher и consumerСредняя, зависит от конфигурации persistence
СкоростьДесятки тысяч сообщений в секундуСотни тысяч операций в секунду
МаршрутизацияExchange, routing keys, bindingsПростые списки и pub/sub
Сложность настройкиВышеНиже
СценарииФинансовые транзакции, заказы, интеграцииКеш, счётчики, лёгкие задачи

RabbitMQ выбирают, когда потеря сообщения недопустима. Redis - когда важна скорость и проект уже использует его как кеш. Для большинства задач с Celery RabbitMQ надёжнее, но Redis быстрее разворачивается и проще в эксплуатации.

Установка RabbitMQ

На Ubuntu/Debian:

sudo apt update
sudo apt install rabbitmq-server
sudo systemctl enable rabbitmq-server
sudo systemctl start rabbitmq-server

На macOS через Homebrew:

brew install rabbitmq
brew services start rabbitmq

На Windows используйте официальный установщик с сайта RabbitMQ или Docker. Запуск через Docker:

docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management

Веб-интерфейс управления доступен на порту 15672. Логин и пароль по умолчанию: guest/guest. Проверьте работоспособность командой:

sudo rabbitmqctl status

Установка Redis

На Ubuntu/Debian:

sudo apt update
sudo apt install redis-server
sudo systemctl enable redis-server
sudo systemctl start redis-server

На macOS:

brew install redis
brew services start redis

Запуск через Docker:

docker run -d --name redis -p 6379:6379 redis:7-alpine

Проверка:

redis-cli ping

Ответ PONG означает, что сервер работает. Подробнее о конфигурации Redis для высоконагруженных систем читайте в руководстве по кешированию в высоконагруженных системах.

Настройка Celery: первый проект

Минимальный проект Celery состоит из файла с задачами и конфигурации брокера. Начните с установки.

Установка Celery и создание приложения

pip install celery

Создайте файл tasks.py:

from celery import Celery

app = Celery(
    'tasks',
    broker='amqp://guest:guest@localhost:5672//',
    backend='rpc://'
)

Параметр broker указывает адрес RabbitMQ. Для Redis используйте redis://localhost:6379/0. Параметр backend определяет, где хранить результаты задач. RPC backend подходит для простых случаев, для production используйте Redis или базу данных.

Определение и вызов первой задачи

@app.task
def add(x, y):
    return x + y

Вызов задачи из Python:

result = add.delay(4, 6)
print(result.id)
print(result.get(timeout=10))

Метод delay() - сокращение для apply_async(). Он отправляет задачу в очередь и возвращает объект AsyncResult. Метод get() блокирует выполнение до получения результата.

Запуск воркера и проверка выполнения

Запустите воркер в отдельном терминале:

celery -A tasks worker --loglevel=info

В логах появится информация о подключении к брокеру и готовности воркера. После вызова add.delay(4, 6) воркер выведет сообщение о получении и выполнении задачи. Результат 10 будет доступен через result.get().

Организация очередей и маршрутизация задач

Одна очередь подходит для простых проектов. Когда задач становится много, их разделяют по приоритетам и типам. Это позволяет запускать отдельные воркеры для разных очередей и контролировать нагрузку.

Создание и назначение очередей

Укажите очередь при вызове задачи:

add.apply_async(args=[4, 6], queue='high_priority')

Воркер, который обрабатывает эту очередь:

celery -A tasks worker --loglevel=info -Q high_priority

Если воркер не подписан на очередь, задача останется в ней до появления подходящего воркера.

Маршрутизация задач по шаблонам

Автоматическая маршрутизация настраивается через task_routes в конфигурации:

app.conf.task_routes = {
    'tasks.add': {'queue': 'math'},
    'tasks.send_email': {'queue': 'notifications'},
}

Теперь при вызове add.delay() задача автоматически попадёт в очередь math. Это удобно для разделения ресурсоёмких и быстрых операций. Отдельные воркеры для каждой очереди позволяют масштабировать только нужные компоненты.

Отложенные и периодические задачи

Celery поддерживает выполнение задач в будущем и по расписанию.

Отложенное выполнение с countdown и ETA

Параметр countdown задаёт задержку в секундах:

send_email.apply_async(args=[user_id], countdown=300)

Задача выполнится через 5 минут. Параметр eta принимает конкретное время:

from datetime import datetime, timedelta

eta = datetime.utcnow() + timedelta(hours=1)
send_email.apply_async(args=[user_id], eta=eta)

ETA использует UTC. При передаче времени в другом часовом поясе Celery может выполнить задачу неверно.

Периодические задачи с Celery Beat

Celery Beat - планировщик, который отправляет задачи по расписанию. Запустите его отдельным процессом:

celery -A tasks beat --loglevel=info

Расписание задаётся в конфигурации:

from celery.schedules import crontab

app.conf.beat_schedule = {
    'daily-report': {
        'task': 'tasks.generate_report',
        'schedule': crontab(hour=3, minute=0),
    },
    'every-30-seconds': {
        'task': 'tasks.cleanup',
        'schedule': 30.0,
    },
}

Задача generate_report будет запускаться каждый день в 03:00. Задача cleanup - каждые 30 секунд. Для production Beat запускают как отдельный сервис, часто вместе с воркером в одном процессе через celery -A tasks worker -B, но для отказоустойчивости лучше разделять.

Интеграция Celery с веб-фреймворками

Веб-приложения - основной потребитель Celery. Длительные операции выносятся из обработчиков запросов в фоновые задачи.

Интеграция с Django

Установите пакет для хранения результатов:

pip install django-celery-results

В settings.py добавьте:

INSTALLED_APPS = [
    'django_celery_results',
]

CELERY_BROKER_URL = 'redis://localhost:6379/0'
CELERY_RESULT_BACKEND = 'django-db'

Создайте файл celery.py в директории проекта:

import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'project.settings')
app = Celery('project')
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()

Задача в tasks.py приложения:

from celery import shared_task

@shared_task
def process_upload(file_id):
    # обработка файла
    return True

Вызов из view:

from django.shortcuts import render
from .tasks import process_upload

def upload_view(request):
    file_id = save_file(request.FILES['file'])
    process_upload.delay(file_id)
    return render(request, 'upload_started.html')

Пользователь сразу получает ответ, а обработка файла идёт в фоне.

Интеграция с Flask

Создайте расширение в extensions.py:

from celery import Celery

celery = Celery(
    'app',
    broker='redis://localhost:6379/0',
    backend='redis://localhost:6379/1'
)

В app.py:

from flask import Flask, request, jsonify
from extensions import celery

app = Flask(__name__)

@celery.task
def send_welcome_email(user_email):
    # отправка письма
    return f'Email sent to {user_email}'

@app.route('/register', methods=['POST'])
def register():
    email = request.json['email']
    send_welcome_email.delay(email)
    return jsonify({'status': 'ok'}), 202

Для задач, которым нужен контекст Flask, используйте with app.app_context(): внутри задачи или передавайте нужные данные аргументами. Запуск воркера: celery -A app.celery worker --loglevel=info.

Мониторинг и управление задачами

Flower - веб-интерфейс для мониторинга Celery. Он показывает активные, завершённые и упавшие задачи, графики производительности и позволяет управлять воркерами.

Использование Flower для мониторинга

Установка:

pip install flower

Запуск:

celery -A tasks flower

Интерфейс доступен на порту 5555. Основные возможности:

  • Список активных, завершённых и отменённых задач
  • Графики количества задач по времени
  • Просмотр аргументов и результатов каждой задачи
  • Управление воркерами: перезапуск, остановка, изменение пула
  • Просмотр очередей и количества задач в каждой

Для production Flower запускают с авторизацией через --basic-auth=user:password и за reverse proxy с HTTPS.

Обработка ошибок и повторные попытки

Задачи падают: сеть недоступна, API вернул ошибку, файл не найден. Celery предоставляет механизм автоматических повторных попыток.

Автоматические повторные попытки

@app.task(bind=True, max_retries=5, retry_backoff=True, retry_jitter=True)
def fetch_data(self, url):
    try:
        response = requests.get(url, timeout=10)
        response.raise_for_status()
        return response.json()
    except requests.RequestException as exc:
        raise self.retry(exc=exc, countdown=60)

Параметры:

  • max_retries - максимальное количество попыток
  • retry_backoff=True - экспоненциальная задержка между попытками
  • retry_jitter=True - случайный разброс задержки для избежания одновременных повторов

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

Обработка исключений и логирование

Логирование внутри задач:

import logging

logger = logging.getLogger(__name__)

@app.task
def process_order(order_id):
    try:
        order = get_order(order_id)
        charge(order)
    except OrderNotFound:
        logger.error('Order %s not found', order_id)
        return False
    except PaymentError as exc:
        logger.exception('Payment failed for order %s', order_id)
        raise
    return True

Для централизованного сбора ошибок подключите Sentry. Пакет sentry-sdk автоматически перехватывает исключения из Celery и отправляет их в Sentry с контекстом задачи. Это ускоряет отладку в production.

Заключение: лучшие практики и дальнейшие шаги

Ключевые моменты этого руководства:

  • RabbitMQ для надёжной доставки, Redis для скорости и простоты
  • Отдельные очереди для разных типов задач упрощают масштабирование
  • Celery Beat для периодических задач, Flower для мониторинга
  • Повторные попытки с экспоненциальной задержкой для устойчивости к сбоям

Для production соблюдайте правила: храните конфигурацию в переменных окружения, запускайте воркеры под супервизором или systemd, настройте мониторинг очередей и алерты. Не храните секреты в коде. Регулярно обновляйте Celery и брокеры.

Дополнительные материалы: автоматизация инфраструктуры для DevOps содержит готовые скрипты для развёртывания сервисов. Для размещения воркеров и брокеров подойдёт облачная инфраструктура Timeweb Cloud с гибким масштабированием ресурсов.

Поделиться:
Сохранить гайд? В закладки браузера