Главный вывод: стандартное решение для асинхронной обработки задач сжатия — очередь 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 мс | Автоповтор + DLQ | 100–100 000 файлов |
| Kubernetes + очередь | Очень высокий (эластичное масштабирование HPA) | <100 мс | Полная отказоустойчивость | Более 100 000 файлов |
Сетевая версия SmartSlim использует архитектуру FastAPI + Celery + Redis + MinIO: на одной машине стабильно работают 12 параллельных задач, а при развёртывании в Kubernetes HPA автоматически масштабирует число реплик от 3 до 10. Эта архитектура уже обеспечивает потребности нескольких корпоративных клиентов с десятками тысяч сжатий в сутки.
2. Подробное описание архитектуры очереди задач Celery
Очередь задач Celery состоит из четырёх ключевых ролей: Producer (производитель), Broker (брокер сообщений), Worker (рабочий процесс) и Backend (хранилище результатов). Чтобы правильно настроить очередь сжатия, важно понимать обязанности каждого слоя.
| Компонент | Технология | Обязанности | Ключевые настройки |
|---|---|---|---|
| Producer | FastAPI | Принимает HTTP-запросы, формирует задачи, отправляет в Broker | task.apply_async(queue=...) |
| Broker | Redis 7.x | Временно хранит сообщения задач, поддерживает приоритетные очереди | broker_url, visibility_timeout |
| Worker | Celery 5.x | Забирает задачи, вызывает движок сжатия на Rust | concurrency, prefork pool |
| Backend | Redis | Хранит состояние и результаты задач | result_backend, result_expires |
| Слой хранения | MinIO | Хранит исходные и сжатые файлы | S3-совместимый протокол, поблочная загрузка |
1. Таблица ключевых параметров Celery
Конфигурация Celery напрямую определяет пропускную способность и стабильность очереди. В таблице ниже — рекомендуемые параметры для сценария сжатия, проверенные на продакшн-среде SmartSlim.
| Параметр | Рекомендуемое значение | Описание |
|---|---|---|
| broker_url | redis://:password@redis:6379/0 | Redis как брокер сообщений, отдельная БД во избежание конфликтов |
| result_backend | redis://:password@redis:6379/1 | Отдельная БД для хранения результатов, изолирована от Broker |
| task_serializer | json | Сериализация в JSON, кросс-языковая совместимость |
| result_serializer | json | Результаты также в JSON |
| accept_content | ['json'] | Принимается только JSON, дополнительная защита |
| timezone | Asia/Shanghai | Единый часовой пояс |
| task_acks_late | True | ACK только после завершения задачи, не теряем задачи при сбое |
| worker_prefetch_multiplier | 1 | Worker берёт только 1 задачу — избегаем голодания долгих задач |
| task_time_limit | 600 | Жёсткий таймаут задачи 600 секунд |
| task_soft_time_limit | 540 | Мягкий таймаут 540 секунд, инициирует SoftTimeLimitExceeded |
| task_reject_on_worker_lost | True | При аварийном завершении Worker задача возвращается в очередь |
| result_expires | 86400 | Результаты автоматически удаляются через 24 часа |
2. Проектирование приоритетных очередей
Задачи сжатия различаются по важности: запросы пользователей в реальном времени требуют быстрого отклика, а задачи ночного архивирования могут выполняться медленно. С помощью приоритетных очередей Redis можно реализовать дифференцированное планирование — Worker в первую очередь забирают высокоприоритетные задачи.
| Имя очереди | Приоритет | Правила маршрутизации | Типичные задачи |
|---|---|---|---|
| compression_high | 9 (наивысший) | Пользовательские запросы в реальном времени | Сжатие одиночного файла по запросу |
| compression_normal | 5 (по умолчанию) | Пакетные задачи | Пакетное сжатие 100–500 файлов |
| compression_low | 1 (низший) | Архивирование по расписанию | Ночное полное архивирование |
| 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 подзадач | Баланс накладных расходов и стоимости повтора |
| Параллелизм Worker | 12 | Режим prefork | 12-ядерный 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 для разных типов.
| Тип ошибки | Типичные исключения | Стратегия | Число повторов | Конечная точка |
|---|---|---|---|---|
| Временная — IO | ConnectionError, TimeoutError | Экспоненциальный откат | 3 | Успех или DLQ |
| Временная — ресурсы | MemoryError, OOMKilled | Длительный откат + понижение | 2 | Успех или DLQ |
| Детерминированная — файл | FileCorrupted, ParseError | Без повтора, сразу DLQ | 0 | dlq_queue |
| Детерминированная — формат | UnsupportedFormat | Без повтора, сразу DLQ | 0 | dlq_queue |
| Детерминированная — доступ | PermissionDenied | Без повтора, оповещение | 0 | dlq_queue + оповещение |
Для повторов используются autoretry_for и retry_backoff в Celery: начальный откат 60 с, максимум 600 с, добавляется случайный джиттер для предотвращения «шторма». Задачи из DLQ периодически сканируются отдельной задачей мониторинга, которая отправляет оповещения в WeChat Work / DingTalk для оперативной обработки.
| Метрика | Способ сбора | Порог оповещения | Действие |
|---|---|---|---|
| Длина очереди | Redis LLEN | > 500 | Запустить масштабирование Worker |
| Доля сбоев задач | Celery events | > 5% | Анализ логов + приостановка приёма |
| Длина DLQ | Redis LLEN dlq | > 10 | Оповещение в WeChat Work |
| Число живых Worker | Celery inspect | < 10 | Автоперезапуск Worker |
| Среднее время задачи | Мониторинг Flower | > 30 с | Проверить крупные файлы + понизить уровень |
| Использование CPU | node_exporter | > 95% | Ограничение + масштабирование |
Полное описание вызова Compression API см. в статьеРуководство по вызову Compression API: дизайн REST-интерфейса。
4. Рекомендации по конфигурации очереди для разных сценариев
Разные бизнес-сценарии предъявляют разные требования к пропускной способности, задержке и надёжности, поэтому конфигурация очереди должна быть дифференцированной. В таблице ниже — рекомендации для типовых случаев.
| Сценарий | Число Worker | Размер фрагмента | Приоритетная очередь | Стратегия повторов |
|---|---|---|---|---|
| Личное сжатие по запросу | 2 | Без фрагментации | high | Быстрый повтор 3 раза |
| Корпоративное пакетное архивирование | 12 | 100 файлов/фрагмент | normal/low | Экспоненциальный откат 3 раза |
| Обработка конфиденциальных документов госструктур | 4 | 50 файлов/фрагмент | high | Строгие повторы + аудит |
| Изображения на e-commerce платформе | 16 | 200 файлов/фрагмент | normal | Быстрый повтор 2 раза |
| Ночное архивирование по расписанию | 8 | 500 файлов/фрагмент | 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 категориях; локальное сжатие без передачи данных за периметр.