Задачи в фоновом режиме с Celery
Если ваша программа имеет задачу с длительным выполнением, например, обработку загруженных данных или отправку электронных писем, вы не хотите ждать её завершения во время запроса. Вместо этого используйте очередь задач, чтобы передать необходимые данные другому процессу, который выполнит задачу в фоновом режиме, в то время как запрос вернёт результат немедленно.
Celery — это мощная очередь задач, которая может использоваться как для простых задач в фоновом режиме, так и для сложных многоэтапных программ и расписаний. Это руководство покажет вам, как настроить Celery с использованием Flask. Прочитайте руководство Celery Первые шаги с Celery, чтобы узнать, как использовать Celery.
Репозиторий Flask содержит пример, основанный на информации на этой странице, который также показывает, как использовать JavaScript для отправки задач и проверки прогресса и результатов.
Установка
Установите Celery из PyPI, например, с помощью pip:
$ pip install celery
Интеграция Celery с Flask
Вы можете использовать Celery без какой-либо интеграции с Flask, но удобно настроить его через конфигурацию Flask и позволить задачам получить доступ к приложению Flask.
Celery использует аналогичные концепции Flask, с Celery объектом приложения, который имеет конфигурацию и регистрирует задачи. При создании приложения Flask используйте следующий код для создания и настройки приложения Celery.
from celery import Celery, Task
def celery_init_app(app: Flask) -> Celery:
class FlaskTask(Task):
def __call__(self, *args: object, **kwargs: object) -> object:
with app.app_context():
return self.run(*args, **kwargs)
celery_app = Celery(app.name, task_cls=FlaskTask)
celery_app.config_from_object(app.config["CELERY"])
celery_app.set_default()
app.extensions["celery"] = celery_app
return celery_app
Это создаёт и возвращает Celery объект приложения. Конфигурация Celery конфигурация берётся из ключа CELERY в конфигурации Flask. Приложение Celery устанавливается в качестве стандартного, чтобы оно отображалось во время каждого запроса. Класс Task автоматически выполняет функции задач с активным контекстом приложения Flask, поэтому доступны такие сервисы, как ваши подключения к базе данных.
Вот базовый example.py пример, который настраивает Celery для использования Redis для связи. Мы включаем бэкенд для результатов, но по умолчанию игнорируем результаты. Это позволяет нам сохранять результаты только для задач, где результат важен.
from flask import Flask
app = Flask(__name__)
app.config.from_mapping(
CELERY=dict(
broker_url="redis://localhost",
result_backend="redis://localhost",
task_ignore_result=True,
),
)
celery_app = celery_init_app(app)
Направьте команду celery worker на этот файл, и она найдёт объект celery_app.
$ celery -A example worker --loglevel INFO
Вы также можете запустить команду celery beat для запуска задач по расписанию. Для получения дополнительной информации о задании расписаний см. документацию Celery.
$ celery -A example beat --loglevel INFO
Фабрика приложения
При использовании шаблона фабрики приложения Flask вызовите функцию celery_init_app внутри фабрики. Она устанавливает app.extensions["celery"] в объект приложения Celery, который можно использовать для получения приложения Celery из приложения Flask, возвращаемого фабрикой.
def create_app() -> Flask:
app = Flask(__name__)
app.config.from_mapping(
CELERY=dict(
broker_url="redis://localhost",
result_backend="redis://localhost",
task_ignore_result=True,
),
)
app.config.from_prefixed_env()
celery_init_app(app)
return app
Чтобы использовать команды celery, Celery нуждается в объекте приложения, но он больше не доступен напрямую. Создайте файл make_celery.py, который вызывает фабрику приложения Flask и получает приложение Celery из возвращаемого приложения Flask.
from example import create_app flask_app = create_app() celery_app = flask_app.extensions["celery"]
Направьте команду celery на этот файл.
$ celery -A make_celery worker --loglevel INFO $ celery -A make_celery beat --loglevel INFO
Определение задач
Использование @celery_app.task для декорирования функций задач требует доступа к объекту celery_app, который не будет доступен при использовании шаблона фабрики. Это также означает, что декорированные задачи привязаны к конкретным экземплярам приложения Flask и Celery, что может быть проблемой во время тестирования, если вы изменяете конфигурацию для теста.
Вместо этого используйте декоратор Celery @shared_task. Это создаёт объекты задач, которые будут получать доступ к тому, что является «текущим приложением», что аналогично концепциям Flask blueprints и контексту приложения. Вот почему мы назвали celery_app.set_default() выше.
Вот пример задачи, которая складывает два числа и возвращает результат.
from celery import shared_task
@shared_task(ignore_result=False)
def add_together(a: int, b: int) -> int:
return a + b
Ранее мы настраивали Celery для игнорирования результатов задач по умолчанию. Поскольку мы хотим узнать возвращаемое значение этой задачи, мы устанавливаем ignore_result=False. С другой стороны, задача, которой не нужно было значение результата, например, отправка письма, этого не сделала бы.
Вызов задач
Декорированная функция становится объектом задачи с методами для вызова её в фоновом режиме. Самый простой способ — использовать метод delay(*args, **kwargs). Для получения дополнительных методов см. документацию Celery.
Для выполнения задачи должен быть запущен Celery worker. Запуск worker показан в предыдущих разделах.
from flask import request
@app.post("/add")
def start_add() -> dict[str, object]:
a = request.form.get("a", type=int)
b = request.form.get("b", type=int)
result = add_together.delay(a, b)
return {"result_id": result.id}
Маршрут не получает результат задачи немедленно. Это бы нарушило цель, блокируя ответ. Вместо этого мы возвращаем идентификатор результата выполняемой задачи, который мы можем использовать позже для получения результата.
Получение результатов
Для получения результата задачи, которую мы запустили выше, мы добавим ещё один маршрут, который принимает идентификатор результата, возвращённый ранее. Мы возвращаем, завершена ли задача (готов ли результат), завершилась ли она успешно и каково было возвращаемое значение (или ошибка), если она завершена.
from celery.result import AsyncResult
@app.get("/result/<id>")
def task_result(id: str) -> dict[str, object]:
result = AsyncResult(id)
return {
"ready": result.ready(),
"successful": result.successful(),
"value": result.result if result.ready() else None,
}
Теперь вы можете запустить задачу с помощью первого маршрута, затем проверить результат с помощью второго маршрута. Это предотвращает блокировку рабочих потоков Flask, ожидающих завершения задач.
Репозиторий Flask содержит пример использования JavaScript для отправки задач и проверки прогресса и результатов.
Передача данных задачам
Задача «add» выше принимала два целых числа в качестве аргументов. Для передачи аргументов задачам Celery необходимо сериализовать их в формат, который он может передать другим процессам. Поэтому передача сложных объектов не рекомендуется. Например, было бы невозможно передать объект SQLAlchemy model, так как этот объект, скорее всего, не сериализуем и привязан к сессии, которая его запросила.
Передавайте минимальный объём данных, необходимых для извлечения или регенерации любых сложных данных в пределах задачи. Рассмотрим задачу, которая выполняется, когда вошедший пользователь запрашивает архив своих данных. Flask запрос знает вошедшего пользователя и имеет объект пользователя, полученный из базы данных. Он получил его, запросив базу данных для данного идентификатора, поэтому задача может сделать то же самое. Передайте идентификатор пользователя, а не сам объект пользователя.
@shared_task
def generate_user_archive(user_id: str) -> None:
user = db.session.get(User, user_id)
...
generate_user_archive.delay(current_user.id)
© 2010 Pallets
Licensed under the BSD 3-clause License.
https://flask.palletsprojects.com/en/stable/patterns/celery/