Distributed Processing System (DP System) — это масштабируемая микросервисная платформа для асинхронной обработки и сжатия видео. Проект спроектирован для работы с файлами большого объема, обеспечивая надежность доставки задач и мониторинг прогресса в реальном времени.
- Language: Go (Golang) 1.24
- API Hero: Gin Gonic (REST API), gRPC (внутреннее взаимодействие между сервисами)
- Database: MySQL 8.0 (хранение метаданных задач)
- Caching & Idempotency: Redis (предотвращение дублирования обработки и Pub/Sub для статусов)
- Messaging: RabbitMQ (отказоустойчивая очередь задач с использованием DLX и политик ретраев)
- Storage: MinIO (S3-совместимое объектное хранилище)
- Processing: FFmpeg (обработка видео через системные вызовы с парсингом прогресса в реальном времени)
- Infrastructure: Docker & Docker Compose (контейнеризация всех компонентов)
- Frontend: React (Vite, TypeScript) с поддержкой чанковой загрузки файлов
Проект демонстрирует решение ряда сложных инженерных задач:
- High-Performance Chunked Upload: Реализована поблочная (chunked) загрузка файлов через
io.Reader. Это позволяет пользователям загружать файлы объемом 3ГБ+ без риска переполнения оперативной памяти сервера. - Resilient Messaging (DLX/Retry): Очередь RabbitMQ настроена с использованием Dead Letter Exchange и Retry Queue. Если воркер не может обработать задачу, она автоматически уходит в очередь ожидания («засыпает» на время через TTL) и пробуется снова, что гарантирует сохранность данных.
- Strict Idempotency: Благодаря интеграции с Redis, система гарантирует, что одно и то же сообщение не будет обработано воркером дважды (защита от дубликатов в RabbitMQ).
- Real-time Progress Tracking: Воркер парсит
out_time_msиз вывода FFmpeg, вычисляет процент выполнения с учетом общей длительности (черезffprobe) и транслирует его клиенту через WebSocket. - Clean Architecture: Код разделен на четкие слои (Domain, Service, Repository, Infrastructure), что упрощает тестирование и масштабирование. Сервисы общаются по gRPC, обеспечивая низкую задержку и строгую типизацию контрактов.
- Загрузка: Клиент инициирует загрузку, разбивает файл на чанки по 5МБ и отправляет их в API.
- Хранение: API собирает чанки и стримит финальный объект в MinIO (Upload Bucket). Метаданные сохраняются в базе данных.
- Очередь: API публикует сообщение о задаче в RabbitMQ.
- Обработка: Воркер забирает задачу, скачивает исходный файл из MinIO во временную папку
/tmpи запускает FFmpeg с выбранным действием (сжатие или конвертация). - Прогресс: Во время работы воркер каждые 1% прогресса или раз в секунду отправляет обновление в Redis Pub/Sub. API слушает этот канал и пересылает данные клиенту по WebSocket.
- Финализация: Результат загружается в Result Bucket, статус задачи в БД обновляется через gRPC-вызов к API, и клиенту отправляется временная presigned-ссылка на скачивание.
- Docker
- Docker Compose
-
Клонировать репозиторий:
git clone https://github.com/your-repo/dpsystem.git cd dpsystem -
Настроить окружение: Создайте файл
.env(или используйте существующий) с необходимыми доступами к БД, RabbitMQ и MinIO. -
Запустить всё одной командой:
docker-compose up -d --build
После запуска клиент будет доступен по адресу: http://localhost:3000.
В результате аудита кода были выявлены и устранены критические узкие места, что позволяет обрабатывать неограниченное количество видео-задач при добавлении вычислительных мощностей.
Проблема: Изначально GORM использовал настройки по умолчанию без ограничения пула соединений. При высокой нагрузке это приводило к ошибке Too many connections и падению MySQL.
Решение:
- Настроен пул соединений
sql.DBчерез GORM. - Установлен
SetMaxOpenConns(50)для ограничения максимального числа открытых соединений. - Установлен
SetMaxIdleConns(20)для переиспользования простаивающих соединений. - Установлен
SetConnMaxLifetime(30m)для предотвращения использования "протухших" соединений.
Проблема: В реализации Hub был обнаружен потенциальный Distributed Deadlock. Метод SendMessage мог синхронно заблокироваться при попытке записи в небуферизированный канал unregister, удерживая RLock, в то время как Hub.Run ожидал этот лок для обработки новых событий. Также отсутствовала потокобезопасность (mutex) при чтении из канала unregister.
Решение:
- Отправка сообщение в канал
unregisterобернута в отдельнуюgo func(), разрывая цикл блокировки. - Добавлен
h.mu.Lock()/Unlock()в обработчикunregisterвнутриRun, устраняя состояние гонки (concurrent map writes). - Убрана O(N) рассылка в
broadcast, отправка сообщений теперь идет адресно за O(1).
Проблема: Обработка видео происходила последовательно или с риском неконтролируемого запуска процессов ffmpeg, что могло привести к CPU starvation.
Решение:
- Внедрен паттерн Worker Pool с фиксированным числом воркеров (5). Это ограничивает максимальную нагрузку на CPU, предотвращая зависание сервера, но сохраняя параллельную обработку задач.
- Добавлен
recover()в горутины воркеров, чтобы паника при обработке "битого" файла не убивала весь процесс воркера.
Проблема: Лимит MaxMultipartMemory был установлен в 100 МБ. При одновременной загрузке множества файлов это могло привести к OOM (Out Of Memory).
Решение:
- Лимит снижен до 16 МБ. Файлы большего размера автоматически стримятся на диск во временные файлы, освобождая RAM для других операций.
Проблема: При резкой остановке сервиса (например, при деплое новой версии или перезагрузке) активные соединения (HTTP-запросы, WebSocket, обработка сообщений из RabbitMQ) обрывались, что приводило к потере данных и ошибкам у пользователей. Решение:
- Реализован перехват системных сигналов
SIGINTиSIGTERM. - При получении сигнала сервер перестает принимать новые HTTP-запросы, но дает время на завершение уже запущенных обработчиков.
- Закрываются соединения с базой данных, Redis и RabbitMQ.
- Воркеры перестают брать новые задачи из очереди, дожидаются завершения текущей обработки и корректно освобождают ресурсы (закрывают файлы, удаляют временные данные).
Developed as a high-performance distributed system showcase.