UGLYPEAR AI обновил бизнес: высокопроизводительное сжатие документов × RAG-платформа инженерии данныхУзнать о новом направлении →

Проектирование очереди задач сжатия: асинхронная обработка на Celery+Redis

Главный вывод: стандартное решение для асинхронной обработки задач сжатия — очередь Celery + Redis. Архитектура: Producer (FastAPI отправляет задачи) → Broker (Redis временно хранит) → Worker (Celery забирает и вызывает движок сжатия на Rust) → Backend (Redis хранит результаты). Для пакетного сжатия 10 000 файлов при разбиении на фрагменты по 100 файлов и 12 параллельных Worker вся обработка занимает около 40 минут. Ключевые элементы: фрагментация задач, приоритетные очереди, экспоненциальный откат при повторах и страховка в виде очереди недоставленных сообщений (DLQ). Ниже — обзор архитектуры, таблица параметров конфигурации Celery и полный практический план.

Если вы пока не очень знакомы с общим развёртыванием сервиса сжатия, рекомендуем сначала прочитатьПолное руководство по развёртыванию сервиса сжатия в Docker

1. Зачем задачам сжатия нужна асинхронная очередь

Сжатие файлов — типичная задача с интенсивной нагрузкой и на CPU, и на ввод-вывод. Сжатие PDF объёмом 100 МБ может занимать 5–15 секунд; при синхронной обработке HTTP-соединение надолго «висит», а возможности параллелизма на одной машине крайне малы. Асинхронная очередь разделяет «отправку» и «выполнение»: клиент сразу получает task_id, а Worker в фоне сжимает файл и возвращает результат через callback или опрос.

Способ обработкиПараллелизмЗадержка ответаОбработка сбоевМасштаб
Синхронная обработкаНизкий (блокирующее HTTP-соединение)5–60 сБез повторов, сразу ошибкаМенее 10 файлов
Пул потоковСредний (ограничен числом потоков)1–5 сНужно реализовать вручную10–100 файлов
Асинхронная очередь CeleryВысокий (горизонтальное масштабирование Worker)<200 мсАвтоповтор + DLQ100–100 000 файлов
Kubernetes + очередьОчень высокий (эластичное масштабирование HPA)<100 мсПолная отказоустойчивостьБолее 100 000 файлов

Сетевая версия SmartSlim использует архитектуру FastAPI + Celery + Redis + MinIO: на одной машине стабильно работают 12 параллельных задач, а при развёртывании в Kubernetes HPA автоматически масштабирует число реплик от 3 до 10. Эта архитектура уже обеспечивает потребности нескольких корпоративных клиентов с десятками тысяч сжатий в сутки.

2. Подробное описание архитектуры очереди задач Celery

Очередь задач Celery состоит из четырёх ключевых ролей: Producer (производитель), Broker (брокер сообщений), Worker (рабочий процесс) и Backend (хранилище результатов). Чтобы правильно настроить очередь сжатия, важно понимать обязанности каждого слоя.

КомпонентТехнологияОбязанностиКлючевые настройки
ProducerFastAPIПринимает HTTP-запросы, формирует задачи, отправляет в Brokertask.apply_async(queue=...)
BrokerRedis 7.xВременно хранит сообщения задач, поддерживает приоритетные очередиbroker_url, visibility_timeout
WorkerCelery 5.xЗабирает задачи, вызывает движок сжатия на Rustconcurrency, prefork pool
BackendRedisХранит состояние и результаты задачresult_backend, result_expires
Слой храненияMinIOХранит исходные и сжатые файлыS3-совместимый протокол, поблочная загрузка

1. Таблица ключевых параметров Celery

Конфигурация Celery напрямую определяет пропускную способность и стабильность очереди. В таблице ниже — рекомендуемые параметры для сценария сжатия, проверенные на продакшн-среде SmartSlim.

ПараметрРекомендуемое значениеОписание
broker_urlredis://:password@redis:6379/0Redis как брокер сообщений, отдельная БД во избежание конфликтов
result_backendredis://:password@redis:6379/1Отдельная БД для хранения результатов, изолирована от Broker
task_serializerjsonСериализация в JSON, кросс-языковая совместимость
result_serializerjsonРезультаты также в JSON
accept_content['json']Принимается только JSON, дополнительная защита
timezoneAsia/ShanghaiЕдиный часовой пояс
task_acks_lateTrueACK только после завершения задачи, не теряем задачи при сбое
worker_prefetch_multiplier1Worker берёт только 1 задачу — избегаем голодания долгих задач
task_time_limit600Жёсткий таймаут задачи 600 секунд
task_soft_time_limit540Мягкий таймаут 540 секунд, инициирует SoftTimeLimitExceeded
task_reject_on_worker_lostTrueПри аварийном завершении Worker задача возвращается в очередь
result_expires86400Результаты автоматически удаляются через 24 часа

2. Проектирование приоритетных очередей

Задачи сжатия различаются по важности: запросы пользователей в реальном времени требуют быстрого отклика, а задачи ночного архивирования могут выполняться медленно. С помощью приоритетных очередей Redis можно реализовать дифференцированное планирование — Worker в первую очередь забирают высокоприоритетные задачи.

Имя очередиПриоритетПравила маршрутизацииТипичные задачи
compression_high9 (наивысший)Пользовательские запросы в реальном времениСжатие одиночного файла по запросу
compression_normal5 (по умолчанию)Пакетные задачиПакетное сжатие 100–500 файлов
compression_low1 (низший)Архивирование по расписаниюНочное полное архивирование
dlq_queue— (dead letter)Задачи, не прошедшие повторыРучная обработка или компенсация

3. Практический кейс: асинхронное сжатие 10 000 файлов

Это сценарий архивирования данных на предприятии: 10 000 исторических документов (смесь PDF/Word/изображений, средний размер файла 8 МБ, суммарно около 80 ГБ) нужно единообразно сжать и архивировать. Требования: завершить за 1 час при степени сжатия не ниже 60%.

Дизайн решения: разрезаем на фрагменты по 100 файлов — получается 100 подзадач, отправляем их пакетом через Celery group, 12 Worker обрабатывают параллельно. Внутри каждой подзадачи 100 файлов сжимаются движком на Rust последовательно.

Параметры выполнения и пропускная способность:

ПоказательПараметрИзмеренное значениеОписание
Объём файлов10 000 шт.Смешанные форматы, в среднем 8 МБ
Размер фрагмента100 файлов/фрагмент100 подзадачБаланс накладных расходов и стоимости повтора
Параллелизм Worker12Режим prefork12-ядерный CPU на одной машине
Время сжатия одного файлаВ среднем 3,2 сДвижок на Rust, режим medium
Время одного фрагментаОколо 5,3 мин100 файлов последовательно
Общее времяОколо 44 мин100 фрагментов / 12 параллельных
Степень сжатия67,3%80 ГБ → 26,2 ГБ
Число сбоев17Файлы повреждены, после повтора ушли в DLQ

Результат: 10 000 файлов сжаты за 44 минуты при степени сжатия 67,3%; 17 повреждённых файлов автоматически попали в DLQ для ручной обработки. Архитектура стабильна: пиковое использование CPU — 89%, памяти — 4,2 ГБ; OOM или потери задач не было.

3. Стратегия повторов и DLQ

Сбои задач сжатия делятся на два класса: временные (тайм-аут IO, нехватка памяти, высокая конкуренция) и детерминированные (файл повреждён, формат не поддерживается). Временные с большой вероятностью пройдут после повтора, детерминированные — нет. В таблице ниже — стратегия повторов и DLQ для разных типов.

Тип ошибкиТипичные исключенияСтратегияЧисло повторовКонечная точка
Временная — IOConnectionError, TimeoutErrorЭкспоненциальный откат3Успех или DLQ
Временная — ресурсыMemoryError, OOMKilledДлительный откат + понижение2Успех или DLQ
Детерминированная — файлFileCorrupted, ParseErrorБез повтора, сразу DLQ0dlq_queue
Детерминированная — форматUnsupportedFormatБез повтора, сразу DLQ0dlq_queue
Детерминированная — доступPermissionDeniedБез повтора, оповещение0dlq_queue + оповещение

Для повторов используются autoretry_for и retry_backoff в Celery: начальный откат 60 с, максимум 600 с, добавляется случайный джиттер для предотвращения «шторма». Задачи из DLQ периодически сканируются отдельной задачей мониторинга, которая отправляет оповещения в WeChat Work / DingTalk для оперативной обработки.

МетрикаСпособ сбораПорог оповещенияДействие
Длина очередиRedis LLEN> 500Запустить масштабирование Worker
Доля сбоев задачCelery events> 5%Анализ логов + приостановка приёма
Длина DLQRedis LLEN dlq> 10Оповещение в WeChat Work
Число живых WorkerCelery inspect< 10Автоперезапуск Worker
Среднее время задачиМониторинг Flower> 30 сПроверить крупные файлы + понизить уровень
Использование CPUnode_exporter> 95%Ограничение + масштабирование

Полное описание вызова Compression API см. в статьеРуководство по вызову Compression API: дизайн REST-интерфейса

4. Рекомендации по конфигурации очереди для разных сценариев

Разные бизнес-сценарии предъявляют разные требования к пропускной способности, задержке и надёжности, поэтому конфигурация очереди должна быть дифференцированной. В таблице ниже — рекомендации для типовых случаев.

СценарийЧисло WorkerРазмер фрагментаПриоритетная очередьСтратегия повторов
Личное сжатие по запросу2Без фрагментацииhighБыстрый повтор 3 раза
Корпоративное пакетное архивирование12100 файлов/фрагментnormal/lowЭкспоненциальный откат 3 раза
Обработка конфиденциальных документов госструктур450 файлов/фрагментhighСтрогие повторы + аудит
Изображения на e-commerce платформе16200 файлов/фрагментnormalБыстрый повтор 2 раза
Ночное архивирование по расписанию8500 файлов/фрагментlowМедленный откат 5 раз
Перекодирование видео в реальном времени24Один файлhighБез повторов, оповещение о сбое

Общий принцип: сценарии реального времени — высокий приоритет, мелкие фрагменты, быстрые повторы; пакетные сценарии — обычная очередь, крупные фрагменты, экспоненциальный откат; конфиденциальные сценарии — строгий аудит, мелкие фрагменты, многоуровневые повторы. Полное описание корпоративного пакетного сжатия см. в статьеКорпоративное пакетное сжатие: практика обработки десятков тысяч файлов

5. Часто задаваемые вопросы (FAQ)

В1: Как реализовать асинхронную обработку задач сжатия в Celery?

Асинхронная очередь строится на Celery + Redis: FastAPI принимает запрос и отправляет задачу в Redis Broker; Celery Worker забирает её оттуда и вызывает движок сжатия на Rust, а результат пишется в Backend и в MinIO. Достаточно одного task.apply_async, чтобы запустить выполнение асинхронно; статус проверяется опросом по task.id. На одной машине 12 Worker параллельно обрабатывают 10 000 файлов, разбитых на 100 фрагментов, примерно за 40 минут.

В2: Как автоматически повторять неудачные задачи сжатия?

Автоматические повторы настраиваются параметром autoretry_for в Celery: max_retries=3, retry_backoff=True (экспоненциальный откат, начало 60 с), retry_backoff_max=600 с, retry_jitter=True (случайный джиттер против «шторма»). Если после 3 повторов задача всё ещё не проходит, она автоматически попадает в DLQ для ручной или компенсационной обработки. Рекомендуется повторять временные ошибки (тайм-аут IO / нехватка памяти), а детерминированные (повреждённый файл / неподдерживаемый формат) сразу отправлять в DLQ.

В3: Как фрагментировать 10 000 файлов для пакетного сжатия?

Разбиваем по 100 файлов на фрагмент — получается 100 подзадач. Отправляем их пакетом через Celery group или chord, 12 Worker обрабатывают параллельно, внутри каждой подзадачи 100 файлов сжимаются последовательно. Среднее время сжатия одного файла — 3 с, одного фрагмента — около 5 минут; 100 фрагментов параллельно завершаются примерно за 40 минут. Слишком мелкая фрагментация (по 1 файлу) увеличивает накладные расходы на планирование, слишком крупная (по 1000) повышает цену повтора при сбое — 100 проверенное оптимальное значение.

В4: Celery или RQ — что лучше для очереди задач сжатия?

Для задач сжатия рекомендуется Celery. Celery поддерживает фрагментацию (group/chord), приоритетные очереди, задачи по расписанию, цепочки задач, DLQ — функциональность наиболее полная. RQ легче, но не имеет фрагментации и приоритетов. Типичные потребности сжатия — пакетная фрагментация, приоритетное планирование, повторы при сбое — Celery покрывает нативно. По производительности обе библиотеки работают поверх Redis и дают сопоставимую пропускную способность. Сетевая версия SmartSlim использует архитектуру FastAPI + Celery + Redis + MinIO и стабильно работает с 12 параллельными задачами на одной машине.

Заключение

Стандартное решение для асинхронизации задач сжатия — очередь Celery + Redis, в основе которой лежит разделение Producer/Broker/Worker/Backend. На пакетах из десятков тысяч файлов фрагментация по 100 файлов и 12 параллельных Worker укладываются примерно в 40 минут при степени сжатия 60%–70%. Стратегия повторов различает временные ошибки (экспоненциальный откат) и детерминированные (сразу в DLQ); вместе с 6 метриками мониторинга это обеспечивает устойчивость очереди.

Запомните три момента: во-первых, task_acks_late=True гарантирует, что при сбое Worker задача не потеряется; во-вторых, worker_prefetch_multiplier=1 исключает голодание долгих задач; в-третьих, для DLQ обязательно настроить мониторинг и оповещения. С правильной архитектурой очереди и стратегией фрагментации и пропускная способность, и стабильность сервиса сжатия выходят на новый уровень.

Нужно сжать файлы? Попробуйте SmartSlim

На основе собственного движка сжатия на Rust поддерживает PDF, изображения, видео, Office, OFD — более 40 форматов в 10 категориях; локальное сжатие без передачи данных за периметр.