Привет, Хабр! Меня зовут Александр Кудряшов. Более тридцати лет я занимаюсь разработкой систем цифровой обработки сигналов, аудио и видео. В последнее время я создаю собственную платформу видеоаналитики и сравниваю её серверные реализации на C++, Python и Rust. Это моя первая статья: в ней хочу поделиться не обзором готового продукта, а одной конкретной инженерной задачей, возникшей при его разработке.
YOLO обрабатывает кадр дольше, чем камера формирует следующий. Если построить видеосервер как последовательность read → motion → YOLO → encode, тяжёлая модель остановит чтение RTSP, очередь начнёт расти, а пользователь увидит не прямой эфир, а прошлое.
В этой статье я разбираю конвейер своего сервера видеоаналитики на C++20. Его основная идея проста: для real‑time‑системы свежесть часто важнее полноты, поэтому между этапами нужны не бесконечные очереди, а ограниченные слоты последнего значения. Покажу, как это сочетается с std::jthread, std::stop_token, OpenCV, ONNX Runtime и защитой от запоздавших результатов.

Постановка задачи
Сервер получает RTSP‑поток с камеры и одновременно должен:
непрерывно читать и декодировать видео;
обнаруживать движение и сопровождать цели;
выполнять полнокадровый YOLO;
дополнительно распознавать небольшие области сопровождаемых объектов;
формировать JPEG‑preview и HLS;
записывать события в PostgreSQL;
по запросу сохранять исходный поток в MP4;
оставаться доступным через REST API.
Камера в моём тесте выдавала примерно 14 кадров/с. Один полнокадровый YOLO‑проход на CPU занимал около 1,7 с. За это время приходило больше двадцати новых кадров.
Если складывать их в обычную FIFO‑очередь, система гарантированно отстанет от реального времени. Увеличение очереди только откладывает момент, когда память закончится или задержка станет неприемлемой.
Поэтому сначала пришлось сформулировать контракт:
Захват видео не ждёт аналитику. Медленный этап обрабатывает самый свежий доступный кадр, а устаревшие кадры допускается заменять.
Это не универсальное правило. Для промышленного контроля, где нельзя пропустить изделие, потребуется другой контракт. Но для наблюдения в реальном времени чаще важнее показать актуальную сцену, чем последовательно обработать уже устаревшие кадры.
Архитектура процесса
Сетевую часть обслуживает Drogon, а видеоконвейер работает в постоянных потоках:
RTSP / VideoCapture | v слот последнего кадра | v координатор / | \ v v v Motion Full Tracking + Track YOLO YOLO 256 \ | / v v v сбор результатов | +----> JPEG / HLS +----> PostgreSQL Drogon HTTP pool ----> REST-команды и состояние FFmpeg processes ----> HLS и запись MP4
Внутри VideoPipeline постоянно живут шесть прикладных потоков:
Поток | Ответственность |
|---|---|
capture | единолично владеет |
coordinator | раздаёт задания и собирает результаты |
motion | выполняет Motion и Tracking |
full YOLO | выполняет полнокадровый inference |
tracking YOLO | распознаёт области активных целей |
display | с заданным темпом готовит JPEG и кадры для HLS |
Кроме них Drogon использует собственный HTTP pool, а два процесса FFmpeg кодируют HLS и MP4.
Потоки создаются один раз. Запускать std::async или новый std::thread для каждого кадра здесь бессмысленно: это добавит накладные расходы и усложнит владение ONNX Runtime Session.
Упрощённо запуск выглядит так:
VideoPipeline::VideoPipeline(...) { m_captureThread = std::jthread( [this](std::stop_token token) { captureRun(token); }); m_motionThread = std::jthread( [this](std::stop_token token) { motionRun(token); }); m_objectThread = std::jthread( [this](std::stop_token token) { objectRun(token); }); m_trackingObjectThread = std::jthread( [this](std::stop_token token) { trackingObjectRun(token); }); m_displayThread = std::jthread( [this](std::stop_token token) { displayRun(token); }); m_thread = std::jthread( [this](std::stop_token token) { run(token); }); }
Почему слот, а не очередь
Между захватом и координатором хранится один последний кадр. Новый кадр заменяет предыдущий, если координатор ещё не успел его забрать.
Упрощённая публикация выглядит так:
void publish(cv::Mat frame) { { std::lock_guard lock(m_captureMutex); m_latestCapture = std::move(frame); ++m_captureSequence; } m_captureChanged.notify_one(); }
Получатель запоминает номер уже обработанного кадра и ждёт изменения:
std::uint64_t consumed = 0; while (!stopToken.stop_requested()) { cv::Mat frame; std::uint64_t sequence = 0; { std::unique_lock lock(m_captureMutex); m_captureChanged.wait(lock, stopToken, [&] { return m_done || m_captureSequence != consumed; }); if (m_done || stopToken.stop_requested()) break; frame = m_latestCapture; sequence = m_captureSequence; } if (sequence > consumed + 1) { profiler.add("pipeline.framesReplacedLatest", sequence - consumed - 1); } consumed = sequence; process(frame); }
Это не lock‑free‑структура, но критическая секция короткая: под mutex меняются только ссылка cv::Mat и номер последовательности. Декодирование, YOLO, JPEG и SQL выполняются после освобождения блокировки.
Поверхностная копия cv::Mat не копирует все пиксели, а увеличивает счётчик ссылок на буфер. Полный clone() нужен только там, где кадр будет изменяться рисованием или должен пережить изменение исходного буфера.
Три разных варианта backpressure
В проекте используются разные ограничители в зависимости от назначения данных:
Последнее значение — между захватом и аналитикой. Устаревший кадр можно заменить.
Один pending task — для YOLO. Пока worker занят, новое задание либо отклоняется, либо заменяет ожидающее.
Небольшая очередь отображения — для HLS. Она сглаживает сетевой джиттер, но имеет жёсткий предел.
Главное — ограничение известно заранее. Память не растёт вместе с длительностью работы сервера.
Worker с std::jthread и stop_token
Для каждого аналитического этапа есть std::optional<Task> и std::optional<Result>, а не очередь произвольной длины:
std::mutex m_objectMutex; std::condition_variable_any m_objectChanged; std::optional<ObjectTask> m_objectPending; std::optional<ObjectResult> m_objectResult; bool m_objectBusy{false};
Worker спит, пока не появится задание или запрос остановки:
void VideoPipeline::objectRun(std::stop_token stopToken) { while (!stopToken.stop_requested()) { ObjectTask task; { std::unique_lock lock(m_objectMutex); m_objectChanged.wait(lock, stopToken, [&] { return m_workersDone || m_objectPending.has_value(); }); if (m_workersDone || stopToken.stop_requested()) break; task = std::move(*m_objectPending); m_objectPending.reset(); m_objectBusy = true; } // Медленный ONNX Run выполняется без удержания mutex. ObjectResult result = detect(task); { std::lock_guard lock(m_objectMutex); m_objectResult = std::move(result); m_objectBusy = false; } } }
std::condition_variable_any здесь выбрана из‑за stop‑aware‑перегрузки wait. Одного request_stop() недостаточно: ожидающий поток нужно также разбудить через notify_all().
В деструкторе сначала останавливаются производители кадров и координатор. Только после их join() завершаются аналитические workers — иначе координатор мог бы положить новое задание в уже остановленный worker.
VideoPipeline::~VideoPipeline() { m_done = true; m_captureThread.request_stop(); m_thread.request_stop(); m_displayThread.request_stop(); m_captureChanged.notify_all(); m_displayChanged.notify_all(); m_captureThread.join(); m_thread.join(); m_displayThread.join(); m_workersDone = true; m_motionThread.request_stop(); m_objectThread.request_stop(); m_trackingObjectThread.request_stop(); m_motionChanged.notify_all(); m_objectChanged.notify_all(); m_trackingObjectChanged.notify_all(); m_motionThread.join(); m_objectThread.join(); m_trackingObjectThread.join(); }
Автоматический join() в деструкторе std::jthread полезен, но он не задаёт требуемый порядок остановки зависимых компонентов. Поэтому порядок сделан явным.
Где stop_token не поможет
stop_token — это совместная отмена, а не принудительное прерывание системного вызова. Например, cv::VideoCapture::read() с FFmpeg backend может блокироваться внутри сторонней библиотеки.
Поэтому для RTSP я задаю таймауты открытия и чтения:
capture.set(cv::CAP_PROP_OPEN_TIMEOUT_MSEC, 5000); capture.set(cv::CAP_PROP_READ_TIMEOUT_MSEC, 5000); capture.set(cv::CAP_PROP_BUFFERSIZE, 1);
После возврата read() цикл снова проверяет stopToken. Цена такого решения — штатная остановка иногда ждёт остаток таймаута.
Корутина сама по себе проблему не решит. Если поместить синхронный VideoCapture::read() или Ort::Session::Run() в coroutine, блокирующая функция не станет неблокирующей. Нужен отдельный worker или настоящий асинхронный адаптер.
Как не принять старый результат за новый
Предположим, пользователь переключил камеру, пока YOLO обрабатывал старый кадр. Через секунду worker вернёт корректный результат — но уже для неправильного источника.
Для защиты используются два поколения:
std::atomic<unsigned long> m_streamGeneration{0}; std::atomic<unsigned long> m_analyticsGeneration{0};
streamGenerationменяется при открытии, остановке или переключении камеры;analyticsGenerationменяется при загрузке модели или изменении настроек аналитики.
Номера записываются в каждое задание и результат:
struct ObjectTask { cv::Mat frame; Json::Value config; std::filesystem::path modelPath; unsigned long streamGeneration; unsigned long analyticsGeneration; };
Координатор применяет результат только при совпадении обоих поколений с текущими. Остальные результаты считаются устаревшими и отдельно учитываются профилировщиком.
Почему поколений два? Изменение параметров YOLO не должно перезапускать RTSP‑сеанс и заставлять ждать новый ключевой кадр камеры. Источник видео и конфигурация аналитики имеют разные жизненные циклы.
Счётчики важнее одного FPS
Средний FPS не объясняет, где исчез кадр. Поэтому я считаю движение каждого задания по конвейеру:
capture.framesOffered capture.framesRead capture.framesFailed pipeline.framesReplacedLatest motion.framesOffered motion.tasksAccepted motion.tasksCompleted motion.tasksSkippedBusy yolo.full2304.tasksAccepted yolo.full2304.tasksCompleted yolo.full2304.tasksSkippedBusy yolo.full2304.resultsDiscardedStale
Для каждого тяжёлого этапа также собираются mean, p50, p95, p99 и max. Это позволяет проверить баланс:
предложено = принято + пропущено busy принято = завершено + выполняется + ошибка
Пропущенные кадры в real‑time‑конвейере допустимы. Необъяснимые кадры — нет.
Что показал ночной тест
Сервер был собран clang-cl 21.1.8 в Release и работал с реальной RTSP‑камерой 2304×1296. Были включены Motion, Tracking, два контура YOLO, PostgreSQL и HLS.
Продолжительность теста составила 16 ч 9 мин 33 с.
Захват и Motion
Метрика | Значение |
|---|---|
Успешно прочитано кадров | 828 302 |
Средняя скорость захвата | 14,2385 кадра/с |
Ошибки чтения | 0 |
Повреждённые кадры | 0 |
Переподключения RTSP | 0 |
Motion обработал | 799 135 кадров |
Motion пропустил как busy | 3 754 кадра |
Motion и Tracking обработали 96,48% прочитанных кадров. Остальные кадры были заменены в слоте последнего значения или не приняты занятым Motion worker.
Полнокадровый YOLO 2304
Метрика | Значение |
|---|---|
Предложено заданий | 802 889 |
Принято | 33 194 |
Завершено | 33 193 |
Пропущено busy | 769 695 |
Ошибки | 0 |
Устаревшие результаты | 0 |
p50 | 1 713,20 мс |
p95 | 1 919,29 мс |
На первый взгляд 4,13% принятых заданий выглядят плохо. Но это ожидаемое поведение: inference занимает около 1,7 с, камера за это время выдаёт десятки кадров, а очередь намеренно не накапливается. YOLO всегда получает свежую работу после завершения предыдущей.
Tracking YOLO 256
Метрика | Значение |
|---|---|
Принято заданий | 59 746 |
Завершено | 58 665 |
Заменено более свежим | 1 081 |
Ошибки | 0 |
p50 | 23,32 мс |
p95 | 61,19 мс |
Здесь баланс сошёлся точно:
59 746 accepted = 58 665 completed + 1 081 replaced
Последний снимок памяти показал 1342 МиБ Working Set, пик — 1409 МиБ. Один тест не доказывает отсутствие утечек, но за 16 часов рост памяти не привёл к остановке процесса или потока.
Что в этой схеме оказалось принципиальным
1. Один владелец блокирующего ресурса
Только capture‑поток вызывает open, read и release у cv::VideoCapture. REST‑команды меняют желаемое состояние и поколение, но не трогают декодер напрямую.
2. Тяжёлая работа вне mutex
Под блокировкой копируются ссылки, параметры и номера поколений. OpenCV, ONNX Runtime, JPEG, FFmpeg и PostgreSQL работают после освобождения mutex.
3. Ограничение памяти является частью архитектуры
Размер каждого буфера выбран явно. Если потребитель медленнее производителя, система заранее знает, что заменить или отбросить.
4. Отмена требует проектирования
jthread упрощает владение потоком, но не отменяет необходимость определить порядок остановки, разбудить условные переменные и учесть блокирующие вызовы библиотек.
5. Результат должен нести контекст
Без поколения камеры и конфигурации асинхронно завершившийся результат невозможно безопасно связать с текущим состоянием системы.
Когда такая архитектура не подходит
Слот последнего кадра полезен, когда важна минимальная задержка. Он не подходит, если требуется:
обработать каждый кадр для последующего аудита;
гарантировать обнаружение краткого события между кадрами;
сохранить строгий порядок всех входных данных;
повторить вычисление после сбоя без потери задания.
В таких системах нужны журналируемая очередь, управление скоростью источника, горизонтальное масштабирование или отдельный контур архивной обработки. Цена полноты — дополнительная задержка и память.
Итог
Real‑time‑конвейер — это не конвейер, который успевает обработать всё. Это система с заранее определённым поведением при перегрузке.
В моём случае рабочей оказалась следующая модель:
один владелец VideoCapture + постоянные workers + слоты последнего значения + явные поколения состояния + stop-aware ожидания + проверяемый баланс счётчиков
C++20 не ускорил YOLO сам по себе. Выигрыш дали нативный код, независимые workers и ограниченные буферы. А std::jthread, std::stop_token и RAII помогли выразить время жизни этой архитектуры так, чтобы её можно было остановить и проверить.
В следующей статье я покажу систему с точки зрения пользователя — от подключения RTSP‑камеры до появления событий в истории. На примере браузерного JS Viewer разберу настройку Motion и Tracking, создание областей анализа, параметры двух контуров YOLO, HLS, запись MP4, PTZ и диагностические показатели. Это позволит связать внутреннее устройство конвейера с тем, как его возможности выглядят и управляются в работающем приложении.

