Некоторые операции в приложении могут занимать много времени и мешать обработке других запросов. Чтобы этого избежать, их выносят в фоновые задачи. В статье разберемся, что это такое, где они применяются и как организовать их выполнение с помощью Celery на нескольких практических примерах.
Что такое фоновые задачи
В веб-разработке и обработке данных нередко возникают ситуации, когда выполнение какой-либо операции может занять значительное время. Представьте, что пользователь загрузил большой файл для обработки, или ему нужно отправить сотни электронных писем. Если эти действия будут выполняться синхронно, непосредственно в момент запроса пользователя, это приведет к «зависанию» приложения, ухудшению пользовательского опыта и, возможно, даже к ошибкам.
Или еще пример: пользователь нажимаете кнопку «Купить» в интернет-магазине. Страница «думает» 10 секунд, потом показывает уведомление «Заказ оформлен». Почему так долго? Потому что в этот момент программа делает все подряд: списывает деньги, создает заказ, отправляет письмо, обновляет статистику. А пользователь ждет.
Но можно сделать иначе: программа быстро принимает заказ и говорит «Готово, мы вас уведомим», а все остальное делает потихоньку в фоне. Пользователь доволен — ему не пришлось долго ждать. Вот для этого и нужны фоновые задачи.
Фоновые задачи (background tasks) – это операции, которые выполняются независимо от основного потока выполнения программы. Проще говоря, это задачи, которые откладываются и выполняются в фоновом режиме, не блокируя работу основного приложения. Это позволяет приложению оставаться отзывчивым, обрабатывать другие запросы и предоставлять пользователю обратную связь о статусе выполнения фоновой задачи.
Давайте рассмотрим несколько простых, но показательных примеров, когда фоновые задачи становятся необходимостью:
- Отправка электронных писем. Отправка одного письма – дело секундное. Но если нужно отправить десятки или сотни писем, например, подтверждение регистрации, уведомления о скидках, рассылки. Синхронная отправка каждого письма может занять значительное время. Фоновая задача позволяет быстро отправить пользователю подтверждение отправки, а письма будут отправляться в фоновом режиме, не задерживая ответ.
- Обработка загруженных файлов. Пользователь загрузил видео для конвертации, или изображение для изменения размера. Эти операции могут быть ресурсоемкими. Вместо того чтобы заставлять пользователя ждать, пока файл обрабатывается, вы можете поставить задачу в очередь и уведомить пользователя, когда обработка будет завершена.
- Генерация отчетов. Создание сложного отчета, который требует агрегации данных из нескольких источников, может занять минуты. Пользователь не должен ждать столько времени. Фоновая задача позволяет ему продолжить работу с сайтом, а отчет будет готов позже.
- Синхронизация данных. Периодическая синхронизация данных с внешними сервисами, например, обновление цен из прайс-листа поставщика, или выгрузка данных в CRM-систему. Эти задачи часто выполняются по расписанию и не должны влиять на работу основного приложения.
- Отложенные уведомления. Отправка напоминаний пользователям о предстоящих событиях, встречах или неоплаченных счетах. Эти уведомления могут быть сгенерированы и поставлены в очередь заранее.
- Импорт/экспорт данных. Массовый импорт данных из CSV-файла или экспорт большого объема информации. Это классический пример задачи, которая может занять много времени.
- Выполнение длительных вычислений. Например, если приложение занимается научными расчетами, моделированием или анализом больших данных, эти вычисления могут занимать часы. Фоновые задачи – единственно возможный вариант для таких сценариев.
В Python существует несколько подходов к реализации фоновых задач. Для простых случаев можно использовать стандартные библиотеки, но для более сложных и масштабируемых решений часто прибегают к специализированным инструментам.
Пример с threading и multiprocessing
Встроенные модули threading и multiprocessing позволяют выполнять код параллельно.
- Потоки (threading). Потоки выполняются в рамках одного процесса и разделяют память. Это подходит для I/O-связанных задач, например, сетевые запросы, чтение/запись файлов, где основной ограничивающий фактор – время ожидания. Однако из-за GIL (Global Interpreter Lock) в CPython, потоки не позволяют достичь истинного параллелизма для CPU-связанных задач.
- Процессы (multiprocessing). Процессы являются независимыми друг от друга и имеют собственную память. Они позволяют обойти GIL и достичь настоящего параллелизма, что идеально подходит для CPU-связанных задач. Однако создание и управление процессами требует больше ресурсов, а обмен данными между ними сложнее.
Простой пример с threading:
import threading
import time
def task(name):
for i in range(1, 4):
print(f"{name} шаг {i}")
time.sleep(1)
# запускаем два потока параллельно
t1 = threading.Thread(target=task, args=("A",))
t2 = threading.Thread(target=task, args=("B",))
t1.start()
t2.start()
# дожидаемся завершения
t1.join()
t2.join()
print("Готово!")
Выделяются два независимых потока выполнения t1 и t2. Через args передаются аргументы для функции task. Сама функция передается через target. В функции генерируются числа 1, 2 и 3, представляющие собой шаги. В консоль выводится имя задачи и текущий шаг. Функция sleep дает задержку в 1 секунду. Когда вызывается start, потоки t1 и t2 начинают работать одновременно. Пока один поток засыпает при вызове sleep, Python переключается на выполнение другого потока. Метод join останавливает главный поток программы и ждет, пока t1 и t2 полностью завершат свою работу. Слово «Готово!» печатается строго в самом конце. Такой же пример, но с multiprocessing:
import multiprocessing
import time
def task(name):
for i in range(1, 4):
print(f"{name} шаг {i}")
time.sleep(1)
# защита от бесконечного порождения процессов на Windows
if __name__ == "__main__":
p1 = multiprocessing.Process(target=task, args=("A",))
p2 = multiprocessing.Process(target=task, args=("B",))
p1.start()
p2.start()
p1.join()
p2.join()
print("Готово!")
Весь основной код находится внутри условной конструкции if __name__ == «__main__»:, которая проверяет, как именно был запущен текущий файл Python: как самостоятельная программа или как импортируемый модуль (библиотека) внутри другого скрипта. Эта проверка критически важна для мультипроцессорности в Windows. Без этой проверки новый процесс снова попытался бы создать еще два процесса, те — еще два, и программа бы зависла или упала из-за бесконечной рекурсии. Код внутри этого блока выполняется только в самом первом (главном) процессе.
В Linux эта проверка технически не требуется для работы multiprocessing, однако ее все равно настоятельно рекомендуется использовать. А если попытаться запустить подобный код в macOS без упомянутой проверки, то программа упадет с ошибкой RuntimeError.
Создаем объекты процессов p1 и p2. Также как и в случае с потоками, передаем им функцию и аргументы для нее. Вызов метода start заставляет операционную систему выделить под каждый процесс отдельное ядро CPU и запустить их одновременно.
Процессы A и B теперь работают параллельно. Метод join останавливает программу и ждет, пока процессы p1 и p2 полностью завершат свою работу (выполнят все 3 шага). Сообщение Готово! появится строго после того, как оба процесса завершатся. Без join() эта строка вывелась бы в самом начале, не дожидаясь выполнения шагов. И в том и в другом случае в консоли шаг за шагом будет выведен следующий результат:
шаг 1 B шаг 1 A шаг 2 B шаг 2 A шаг 3 B шаг 3 Готово!
Celery — мощный инструмент для фоновых задач
Для более сложных систем, где требуется надежное выполнение задач, масштабирование, мониторинг и управление очередями, используются специализированные библиотеки. Самой популярной и мощной из них в экосистеме Python является Celery.
Celery – это распределенная система очередей задач, написанная на Python. Она позволяет легко создавать, управлять и масштабировать фоновые задачи. Как работает Celery?
Архитектура Celery состоит из трех основных компонентов:
- Воркеры (Workers). Это процессы, которые слушают очередь задач и выполняют сами задачи. Можно запускать множество воркеров на разных серверах для масштабирования.
- Брокер сообщений (Message Broker). Это «почтовое отделение», куда отправляются задачи. Воркеры забирают задачи из брокера. Самые популярные брокеры для Celery – RabbitMQ и Redis.
- Задачи (Tasks). Это функции Python, которые необходимо выполнить в фоновом режиме. Нужно декорировать обычную функцию Python декоратором @celery_app.task, чтобы сделать ее задачей Celery.
Процесс выполнения задачи в Celery:
- Основное приложение отправляет задачу (вызов декорированной функции) в брокер сообщений.
- Брокер помещает сообщение о задаче в очередь.
- Один из свободных воркеров забирает сообщение из очереди.
- Воркер выполняет код задачи.
- Результат выполнения (если он нужен) может быть сохранен в базе данных или другом хранилище.
Рассмотрим самый простой пример использования Celery. Но сначала надо установить все необходимые компоненты. В первую очередь надо установить сам Celery и Redis:
pip install celery redis
Далее нам понадобится Docker-контейнер, в котором будет крутиться Redis. Сначала нужно установить сам Docker. На этой странице есть подробные инструкции по установке Docker для разных ОС. После установки Docker запускаем его следующей командой:
docker run -d -p 6379:6379 redis:alpine
Теперь можно приниматься за исходники. У нас будет два файла с расширением .py: в первом будет задача, которую надо выполнить, а второй будет вызывать эту задачу. Первый файл tasks.py:
from celery import Celery
import time
app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0')
@app.task
def slow_multiplication(a, b):
time.sleep(5) # притворяемся, что долго работаем
return a * b
Здесь мы инициализируем экземпляр Celery с именем app, для которого указываем имя главного модуля, а также брокер с бэкэндом, в котором будут храниться результаты. Декорируем функцию slow_multiplication специальным декоратором @app.task. Теперь эта функция является задачей Celery. Функция через 5 секунд возвращает произведение двух чисел. Теперь можно запустить воркер:
celery -A tasks worker --loglevel=info
После успешного запуска воркера можно вызывать задачу. Вызов задачи нужно производить в другом терминале, а не в том, в котором запускался воркер. Можно открыть как новую вкладку, так и новое окно терминала. Главное условие — запустить процесс Celery в отдельной сессии командной строки. Второй файл app.py:
from tasks import slow_multiplication result = slow_multiplication.delay(10, 6) # задача улетела в очередь print(result.ready()) # False - еще считает print(result.get()) # 60 - готово! Получили результат вычисления print(result.ready()) # True - посчитал
В первой строчке делаем импорт функции slow_multiplication из первого файла. Как мы уже говорили, эта функция является задачей Celery, поэтому отправляем ее в очередь при помощи метода delay. В итоге функция не запускается напрямую, а отправляется в очередь задач (брокер). Функция начинает выполняться в фоне (в отдельном процессе-воркере), а код сразу же идет дальше, возвращая объект AsyncResult в переменную result.
Первый вызов метода ready для result проверяет, завершилось ли выполнение задачи. Возвращается False, так как воркер еще считает результат. Метод get останавливает выполнение основного кода и «ждет», пока фоновая задача завершится. Как только воркер закончит вычисления, get() вернет результат. Ну, и второй вызов ready теперь возвращает True, подтверждая, что задача успешно выполнена и результат получен. Открываем новое окно терминала или новую вкладку и вызываем задачу:
python3 app.py
Сразу после запуска app.py в консоли увидим такой вывод:
False
А спустя 5 секунд, вывод уже будет такой:
False 60 True
Далее разберем несколько практических примеров применения Celery.
Генерация пароля
Для генерации пароля будем использовать встроенный в стандартную библиотеку Python модуль secrets, способный создавать криптографически стойкие случайные числа и модуль string с коллекциями строковых констант. Оба эти модуля необходимо сначала импортировать:
import secrets import string
Также не забываем импортировать Celery и создать его экземпляр, как это было показано выше:
from celery import Celery
app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0')
Функция для генерации пароля:
@app.task def generate_password(length): chars = string.ascii_letters + string.digits + '!@#$%^&*' return ''.join(secrets.choice(chars) for _ in range(length))
В переменной chars содержится набор символов, состоящий из букв, чисел и некоторых специальных символов. С помощью choice создается криптографически безопасный набор символов. Метод choice выбирает из набора chars случайный символ и этот выбор повторяется до тех пор, пока не будет сконструирован пароль требуемой длины length. Вызов задачи:
from tasks import generate_password
password = generate_password.delay(16)
print('Пароль:', password.get())
Импортируем функцию из tasks. Вызываем задачу generate_password с аргументом 16 (длина пароля). Задача отправляется в очередь на выполнение, а переменная password сохраняет объект AsyncResult. Метод get переводит выполнение кода в блокирующий режим. Программа останавливается и ждет, пока воркер выполнит задачу и вернет сгенерированный пароль. Как только задача завершится, get() вернет саму строку с паролем, а функция print выведет ее на экран. После запуска воркера и вызова задачи, в консоли получим примерно следующее:
Пароль: J1XKcG7Z6%ZuztV#
Получили надежный пароль из 16-и символов.
Парсинг RSS и извлечение заголовков
Для парсинга нам понадобится библиотека feedparser. Устанавливаем ее следующей командой:
pip install feedparser
Далее ее нужно импортировать в самом начале исходника. Функция для парсинга:
@app.task def rss_titles(feed_url): feed = feedparser.parse(feed_url) return [entry.title for entry in feed.entries]
Парсим ленту при помощи parse, а потом генератором извлекаем из нее заголовки и возвращаем их как результат выполнения функции. Для примера извлечем первые 6 заголовков из RSS-ленты Хабра. Вызов задачи:
from tasks import rss_titles
result = rss_titles.delay('https://habr.com/ru/rss/all/')
titles = result.get()
for title in titles[:6]:
print(title)
Celery отправляет задачу в очередь сообщений (брокер). Функция выполняется асинхронно отдельным процессом (воркером), не блокируя основной поток кода. Метод get заставляет код остановиться и дождаться ответа от воркера. Как только воркер закончит скачивать и обрабатывать RSS-ленту, get вернет результат в виде списка заголовков и запишет его в переменную titles. Далее идет цикл for, который проходит по срезу из 6-и заголовков и выводит их в консоль. После вызова задачи получим примерно такой вывод:
Генератор КИТУ или где обитает весь зоопарк упаковочных кодов: КИГУ, КИН, SSCC и АТК От статьи на Хабре до Байконура: как я еду в космос [Перевод] Эффективность использования капитала: метрика, которую игнорирует большинство пользователей DeFi Речевая аналитика на домашнем «чайнике» Еще одна айтишно-заводская задача: отслеживаем историю прокатных валков, чтобы снизить простои Анатомия граблей
Скачивание файла по URL
Чтобы скачать файл по URL нам понадобится библиотека requests, поэтому сначала надо ее установить:
pip install requests
А потом импортировать. Сама функция для загрузки файла выглядит так:
@app.task
def download_file(url, save_path):
response = requests.get(url, timeout=10)
response.raise_for_status()
with open(save_path, 'wb') as f:
f.write(response.content)
return {'path': save_path, 'bytes': len(response.content)}
В функции создается GET-запрос, в котором передается URL и таймаут. С помощью конструкции with open открываем файл для записи по пути, где хотим сохранить загруженный файл. Записываем файл методом write. Вызываем задачу:
from tasks import download_file
result = download_file.delay(
'https://example.com/file.zip',
'/home/user/downloads/file.zip'
)
print(result.get())
После импорта функции, отправляем ее в очередь и начинаем выполнение задачи в фоне. В качестве аргументов передаем функции адрес файла и место, куда его надо сохранить. Методом get останавливаем основную программу и ждем пока файл загрузится. После завершения загрузки, get возвращает результат.
Скачаем, например, пьесу Марии Смагиной «Ангел летит» из библиотеки пьес Александра Чупина, которая находится по этому адресу. Код вызова задачи будет примерно таким:
from tasks import download_file
result = download_file.delay("https://krispen.ru/smagina_m_01.docx", "/home/alex/file.docx")
print(result.get())
Вывод в консоли будет следующим:
{'path': '/home/alex/file.docx', 'bytes': 34321}
А в домашней директории появится файл с именем file.docx.
Подсчет строк в CSV-файле
CSV (Comma-Separated Values) — это простой текстовый формат предназначенный для хранения и передачи табличных данных. Каждая строка в файлах такого формата соответствует строке таблицы, а значения разделяются запятыми или другими символами. Подробнее о CSV можно прочитать в этой статье.
Создадим простенький csv-файл вот с таким содержимым:
Name,Age,City Alex,28,New York Anna,54,Paris Peter,45,London Olga,22,Moscow
Назовем его, например, people.csv и положим в корень домашней директории. Чтобы подсчитать количество строк в этом файле при помощи Celery нам понадобится такая функция:
@app.task def count_csv_rows(filepath): with open(filepath, newline='', encoding='utf-8') as f: reader = csv.reader(f) return sum(1 for _ in reader)
Перед созданием экземпляра Celery надо импортировать встроенный в стандартную библиотеку модуль csv, так как в коде используется метод reader из этого модуля. Открываем файл также при помощи with open, как и в прошлом примере, но уже для чтения, а не для записи. Читаем файл при помощи reader, который разбирает файл на строки и автоматически разбивает их на списки элементов согласно запятым или другим разделителям.
Подсчет ведем при помощи простого генератора. Он проходит по файлу строка за строкой и за каждую найденную строку выдает единицу, а функция sum складывает эти единицы. Код вызова задачи:
from tasks import count_csv_rows
result = count_csv_rows.delay('/home/alex/people.csv')
print('Строк в файле:', result.get())
Вывод в консоли:
Строк в файле: 5
Так как первая строка файла — это обычно шапка таблицы, а нужно только число строк с данными, то от результата следует отнять единицу.
Выполнение задачи по расписанию
Ну и напоследок рассмотрим парочку примеров запуска задач по расписанию. Допустим, нам надо, чтобы какая-то задача выполнялась каждые две минуты. Конечно, можно использовать циклы, но как с этим может помочь Celery? Для таких вещей есть Celery Beat. Это такой планировщик задач в экосистеме Celery, который запускает задачи по расписанию с определенной периодичностью и в определенное время.
Он не выполняет задачи сам, а отправляет их в очередь, откуда их подхватывают воркеры. Пусть нам нужно каждые 2 минуты выводить в консоль какое-то сообщение. Вот так это можно сделать:
from celery import Celery
from datetime import timedelta
app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0')
@app.task
def print_elapsed():
print("Прошло 2 минуты.")
app.conf.beat_schedule = {
'print-every-2-minutes': {
'task': 'tasks.print_elapsed',
'schedule': timedelta(minutes=2), # каждые 2 минуты
},
}
Здесь нам, в принципе, все знакомо, кроме последней части с фигурными скобками. Это и есть расписание. В нем через task указывается функция, которую надо выполнить, а через schedule — периодичность, с которой эта функция должна выполняться. В нашем случае периодичность указываем при помощи класса timedelta из встроенного модуля datetime. Для выполнения задачи сначала нужно запустить воркер уже знакомой командой:
celery -A tasks worker --loglevel=info
А потом в отдельном терминале запустить планировщик:
celery -A tasks beat --loglevel=info
В результате в терминале воркера каждые 2 минуты будет появляться примерно такой вывод:
[2026-09-15 14:01:09,889: INFO/MainProcess] Task tasks.print_elapsed[bde76928-41f7-4a1d-9eb1-91a10cfa8df7] received [2026-09-15 14:01:09,891: WARNING/ForkPoolWorker-4] Прошло 2 минуты [2026-09-15 14:01:09,893: INFO/ForkPoolWorker-4] Task tasks.print_elapsed[bde76928-41f7-4a1d-9eb1-91a10cfa8df7] succeeded in 0.0020716090002679266s: None
А в терминале планировщика такой:
[2026-09-15 14:01:09,883: INFO/MainProcess] Scheduler: Sending due task print-every-2-minutes (tasks.print_elapsed)
Все прекрасно работает. Задача выполняется. Также можно указать определенное время вывода сообщения. Следующий код каждую полночь выводит соответствующее сообщение:
from celery import Celery
from celery.schedules import crontab
app = Celery('tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0')
@app.task
def print_midnight_message():
print("Ровно полночь!")
app.conf.update(timezone='Asia/Irkutsk', enable_utc=False)
app.conf.beat_schedule = {
'midnight-daily': {
'task': 'tasks.print_midnight_message',
'schedule': crontab(minute=0, hour=0), # каждый день в 00:00
},
}
Вместо timedelta здесь используется crontab из модуля celery.schedules. Так как Celery Beat по умолчанию работает в UTC, то сразу после функции для вывода сообщения происходит обновление конфигурации с учетом часового пояса. На последнем этапе в конфигураторе расписания с помощью crontab указывается время запуска задачи. Запускается код так же, как и в прошлом примере и в полночь в консоли должно появиться указанное сообщение:
......WARNING/ForkPoolWorker-4] Ровно полночь!
Для чего может понадобиться Celery Beat? Например, для рассылки уведомлений, напоминаний или писем. Также его можно использовать для автоматического продления подписок. Еще он неплохо справляется с ежедневной генерацией отчетов и выгрузкой данных. Или, если речь идет о регулярной синхронизации с внешними API, то и здесь Celery Beat будет прекрасным помощником.
