Разработка · 6 мин чтения

Telegram бот с очередью задач Celery: зачем нужна и как подключить

Первый бот, который я угробил таймаутом Telegram API, рассылал статусы заказов: цикл по 3000 подписчикам прямо в хендлере aiogram, и через полторы минуты бот переставал отвечать даже на /start, потому что event loop был занят рассылкой. С тех пор на любом проекте, где есть массовые действия - рассылки, интеграция со СДЭК, генерация отчётов, - я сразу закладываю telegram бот с очередью задач celery в архитектуру, а не пристраиваю её потом, когда бот уже подвис у заказчика в проде.

Зачем боту очередь, если aiogram и так асинхронный

Асинхронность aiogram решает конкурентность внутри одного процесса: пока бот ждёт ответ от одного пользователя, он успевает обработать сообщение от другого. Но у неё есть три места, где она не спасает.

Первое - лимиты самого Telegram Bot API: 30 сообщений в секунду на бота в целом и не больше 1 сообщения в секунду в один чат. Если в хендлере просто раскидать asyncio.gather на 3000 send_message, часть уйдёт в FloodWait, а необработанное исключение в цикле может уронить весь обработчик обновления.

Второе - задачи, которые должны пережить рестарт процесса. Деплой, падение контейнера, обновление кода - если рассылка или обработка вебхука сидела в памяти процесса бота, при рестарте она просто теряется. Задача в очереди Celery лежит в Redis или RabbitMQ и подхватится воркером, даже если сам бот в этот момент перезапускается.

Третье - масштабирование. Обработку тяжёлых задач (генерация PDF-отчётов, ресайз изображений, запросы к нестабильному внешнему API) хочется вынести на отдельные процессы или даже отдельный сервер, не трогая процесс, который держит polling или webhook. Celery-воркеры масштабируются горизонтально независимо от самого бота - добавил контейнер, увеличил concurrency, и очередь разгребается быстрее.

Redis или RabbitMQ - какой брокер ставить под Celery

Брокер - это то, куда Celery кладёт задачи и откуда их забирают воркеры. Для бота выбор обычно между Redis и RabbitMQ.

Критерий Redis RabbitMQ
Настройка Один контейнер, работает сразу Требует настройки vhost, exchange, routing key
Гарантия доставки At-least-once, приемлемо для большинства задач Строже: подтверждения, persistent-очереди
Нагрузка на память Низкая, хватает 512 МБ RAM на VPS Выше из-за брокера сообщений и Erlang VM
Когда ставлю 90% ботов: рассылки, уведомления, парсинг Высоконагруженные вебхуки оплаты, где важна каждая транзакция

В большинстве ботов на aiogram я ставлю Redis по простой причине: он уже нужен для хранения состояний FSM (RedisStorage). Разношу по базам - db=0 под FSM, db=1 под брокер Celery, db=2 под backend с результатами задач. Один контейнер, три логические зоны, не нужно поднимать RabbitMQ ради очереди на десяток задач в минуту.

Как подключить Celery к aiogram-боту: минимальная связка

Структура простая: celery_app.py с инициализацией приложения, tasks.py с самими задачами, и bot.py, где хендлеры вызывают задачи через .delay() вместо того, чтобы выполнять тяжёлую работу внутри себя.

# celery_app.py
from celery import Celery

celery_app = Celery(
    "bot_tasks",
    broker="redis://localhost:6379/1",
    backend="redis://localhost:6379/2",
)

celery_app.conf.update(
    task_serializer="json",
    result_serializer="json",
    accept_content=["json"],
    timezone="Europe/Moscow",
    task_acks_late=True,
    worker_prefetch_multiplier=1,
)

# tasks.py
import asyncio
from aiogram import Bot
from celery_app import celery_app

BOT_TOKEN = "..."

@celery_app.task(bind=True, max_retries=3, default_retry_delay=30)
def send_bulk_notification(self, user_ids: list[int], text: str):
    async def _send():
        bot = Bot(token=BOT_TOKEN)
        for chat_id in user_ids:
            try:
                await bot.send_message(chat_id, text)
                await asyncio.sleep(0.05)  # держим лимит 1 msg/sec в чат с запасом
            except Exception as exc:
                raise self.retry(exc=exc)
        await bot.session.close()

    asyncio.run(_send())

В хендлере aiogram вызов задачи выглядит так: send_bulk_notification.delay(user_ids, text). Это не блокирующий вызов - он просто кладёт задачу в Redis и сразу возвращает управление, а сама рассылка выполняется в отдельном процессе воркера.

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

celery -A celery_app worker --loglevel=info --concurrency=4
celery -A celery_app beat --loglevel=info
celery -A celery_app flower --port=5555

Бесплатный материал

🎁 Полезный скрипт в подарок

Подпишитесь на Telegram - пришлю готовый скрипт по этой теме.

Без спама. Отписка в 1 клик.

Практические кейсы: рассылки, вебхуки СДЭК и оплата T‑Bank/WooCommerce

Рассылки - самый частый повод подключить очередь. Вместо цикла в хендлере разбиваю список подписчиков на пачки по 25-30 и отправляю каждую пачку отдельной задачей с небольшим countdown между ними - так проще ловить FloodWait от конкретного чата и ретраить именно его, не завися от судьбы всей рассылки.

Вебхуки СДЭК - второй частый случай. СДЭК ждёт ответ 200 на уведомление о смене статуса заказа в течение нескольких секунд, иначе повторяет запрос. Если внутри хендлера вебхука делать поход в базу, потом в Telegram с уведомлением клиенту, а потом ещё синхронизацию с CRM, легко не уложиться в таймаут и получить дубли уведомлений. Правильный порядок: хендлер вебхука проверяет подпись, кладёт задачу в очередь и сразу отвечает 200, а сама обработка (обновление статуса, уведомление в боте, запись в CRM) идёт в celery-таске с retry на случай, если CRM в этот момент недоступна.

Тот же принцип с оплатой через T‑Bank или синхронизацией заказов с WooCommerce: вебхук о поступлении платежа обязан подтвердиться быстро, а списание остатков на складе, отправка чека и сообщение покупателю в боте переносятся в очередь. Если WooCommerce на секунду отдал 500‑ю ошибку при синхронизации остатка, задача уйдёт на retry через 30-60 секунд, а не потеряется вместе с оплаченным заказом.

Отдельно использую Celery beat для периодических задач - мониторинг цен конкурентов, проверка статусов у СДЭК, выгрузка отчётов раз в сутки. Раньше решал это через системный cron, но когда таких задач становится больше пяти, удобнее держать расписание в одном месте вместе с остальной логикой очереди. Похожие заготовки для парсинга и фоновых уведомлений у меня уже собраны в библиотеке готовых скриптов, если не хочется писать таск с нуля под типовую задачу.

Мониторинг и масштабирование воркеров в проде

Flower - веб-интерфейс для Celery, который показываю заказчикам вместо консольных логов: видно очередь задач, сколько упало с retry, сколько выполняется прямо сейчас, средняя скорость обработки. Поднимаю его отдельным контейнером на порту 5555 и закрываю базовой авторизацией через nginx.

Когда в одном проекте соседствуют быстрые задачи (уведомление в боте) и медленные (обработка вебхука с внешним API, который иногда отвечает 10 секунд), развожу их по отдельным очередям: celery -A celery_app worker -Q notifications --concurrency=2 и отдельный воркер с -Q webhooks --concurrency=4. Так медленный вебхук не держит в ожидании быстрые уведомления пользователям.

Важный момент по вебхукам оплаты и СДЭК - идемпотентность. Retry означает, что одна и та же задача может выполниться дважды, поэтому в таске проверяю по id заказа, не обработан ли он уже, прежде чем что-то менять в базе или отправлять повторное уведомление клиенту. Без этой проверки при третьем retry клиент получит три одинаковых сообщения о статусе заказа - сталкивался с этим на одном из первых проектов с интеграцией СДЭК, пока не добавил проверку через уникальный ключ идемпотентности в Redis.

Автоматизация в мессенджере

Telegram-бот / Mini App

от 30 000 ₽

Подробнее →

Частые вопросы

Обязательно ли ставить Celery, если у бота мало пользователей и нет тяжёлых задач

Нет. Для простого бота без вебхуков и массовых рассылок хватает обычного asyncio.create_task внутри хендлера или встроенной очереди на asyncio.Queue. Celery оправдан, когда задача должна пережить рестарт процесса, нуждается в автоматическом retry или когда обработку нужно разнести по нескольким серверам. Для совсем нетехнических сценариев (простое расписание рассылок без кастомной логики) иногда проще собрать цепочку в n8n, чем поднимать отдельную инфраструктуру под Celery.

Можно ли обойтись без Redis и использовать SQLite как брокер

Технически у Celery есть транспорт через SQLAlchemy, но под конкурентной записью от нескольких воркеров SQLite начинает блокироваться и терять задачи - в проде так не делаю. Redis в качестве брокера занимает минимум ресурсов и спокойно работает на VPS с 512 МБ памяти, поэтому смысла экономить на нём нет.

Как Celery работает, если бот не на polling, а принимает апдейты через вебхук на FastAPI

Так же, как и с polling: хендлер вебхука принимает апдейт от Telegram, быстро отвечает 200 и кладёт тяжёлую часть обработки в celery-таск через .delay(). Брокер и воркеры не зависят от того, как бот получает апдейты - важно только, что сам процесс с FastAPI или aiohttp не блокируется на время выполнения задачи.

Сколько стоит добавить очередь задач в уже работающего бота

Если бот уже написан на aiogram и есть доступ к серверу, доработка сводится к подключению Celery, вынесению рассылок или вебхуков в задачи и настройке воркера рядом с ботом. Такие интеграции с настройкой очереди и деплоем оцениваю от 30 000 ₽ - итоговая сумма зависит от количества сценариев, которые нужно перевести на асинхронную обработку, и от того, нужен ли отдельно RabbitMQ вместо Redis.

Есть задача?

Обсудим в мессенджере

Расскажите, что нужно сделать — отвечу в течение 4 часов в рабочее время. Первая консультация бесплатно.

Самозанятый Калинкин Н. А. · работаю с физлицами и юрлицами

Продолжая пользование настоящим сайтом Вы выражаете своё согласие на обработку Ваших персональных данных (файлов куки) с использованием Yandex.Metrika.
Понятно