Введение в очереди сообщений и Celery
Очередь сообщений решает задачу асинхронной обработки данных. Вместо того чтобы заставлять пользователя ждать завершения длительной операции, приложение отправляет задачу в очередь и сразу отвечает. Брокер хранит задачу, а воркер забирает её и выполняет в фоне. Это снижает задержку ответа и повышает отказоустойчивость системы.
Celery - это распределённая очередь задач для Python. Библиотека берёт на себя управление задачами: отправку, маршрутизацию, повторные попытки и хранение результатов. Celery работает с брокерами сообщений RabbitMQ и Redis, которые отвечают за передачу задач от приложения к воркерам. RabbitMQ обеспечивает надёжную доставку и сложную маршрутизацию, Redis - высокую скорость и простоту настройки.
Типовые сценарии использования Celery: отправка email после регистрации пользователя, обработка изображений, генерация отчётов, синхронизация данных с внешними API. Любая операция, которая занимает больше пары секунд, кандидат на вынос в фоновую задачу. Это разгружает веб-сервер и делает интерфейс отзывчивым.
Базовые принципы работы очередей сообщений разобраны в практическом руководстве по оптимизации логирования в Python. Для выбора между RabbitMQ и Redis полезно изучить алгоритм выбора брокера сообщений.
Установка и настройка брокера сообщений
Первый шаг - установить брокер. От выбора зависит надёжность доставки, скорость и сложность эксплуатации.
Выбор брокера: RabbitMQ или Redis?
| Критерий | RabbitMQ | Redis |
|---|---|---|
| Надёжность доставки | Высокая, подтверждения 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 с гибким масштабированием ресурсов.