PunchNLNG unveils AI tool to cut MRI scan timeThe Jerusalem PostWearable counter‑drone tech enters frontline service across US, Ukraine and IDF unitsBollywood HungamaMahakavya Shri Ramayan Katha producer Prakash Mahobiya alleges ‘negative marketing’; hints at big-budget Ramayana saying, “We never got a chance to reach our audience”ESPNTransfer rumors, news: Man United, Arsenal, Chelsea battle for Freiburg strikerInquirerSara Duterte: Mom prefers Baste for national politicsCNN Türk"Fon" ödemesi hangi formülle olacak?UOLTelevangelista americano Jim Bakker, envolvido em escândalos de fraude e sexo, morre aos 86 anos20 MinutenXena (24): «Sie trauen mir als Frau den Chefposten nicht zu»ZDF heuteEntdecken Sie das ZDF-NachrichtenstudioHet Laatste NieuwsVoormalig hoofd van Duitse inlichtingendienst aangehouden op verdenking van spionage en landverraadWirtualna PolskaPrezes UOKiK: Liczymy na refleksję po stronie Google'aNHK 社会デヴィ夫人 元マネージャーなど暴行の罪で罰金20万円
The Daily Newsstand · Free, Always
Tuesday, October 6, 2026

Как не утопить RTSP в YOLO: real‑time‑конвейер на C++20, Clang и ONNX Runtime

Translate

Привет, Хабр! Меня зовут Александр Кудряшов. Более тридцати лет я занимаюсь разработкой систем цифровой обработки сигналов, аудио и видео. В последнее время я создаю собственную платформу видеоаналитики и сравниваю её серверные реализации на C++, Python и Rust. Это моя первая статья: в ней хочу поделиться не обзором готового продукта, а одной конкретной инженерной задачей, возникшей при его разработке.

YOLO обрабатывает кадр дольше, чем камера формирует следующий. Если построить видеосервер как последовательность read → motion → YOLO → encode, тяжёлая модель остановит чтение RTSP, очередь начнёт расти, а пользователь увидит не прямой эфир, а прошлое.

В этой статье я разбираю конвейер своего сервера видеоаналитики на C++20. Его основная идея проста: для real‑time‑системы свежесть часто важнее полноты, поэтому между этапами нужны не бесконечные очереди, а ограниченные слоты последнего значения. Покажу, как это сочетается с std::jthread, std::stop_token, OpenCV, ONNX Runtime и защитой от запоздавших результатов.

Браузерный Viewer: RTSP-кадр 2304×1296, области анализа, сопровождение целей и результаты YOLO. В строке состояния отображаются FPS сервера и клиента, потери, очередь и число перезапусков RTSP.

Браузерный Viewer: RTSP‑кадр 2304×1296, области анализа, сопровождение целей и результаты YOLO. В строке состояния отображаются FPS сервера и клиента, потери, очередь и число перезапусков RTSP.

Постановка задачи

Сервер получает 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

единолично владеет cv::VideoCapture

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

В проекте используются разные ограничители в зависимости от назначения данных:

  1. Последнее значение — между захватом и аналитикой. Устаревший кадр можно заменить.

  2. Один pending task — для YOLO. Пока worker занят, новое задание либо отклоняется, либо заменяет ожидающее.

  3. Небольшая очередь отображения — для 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 и диагностические показатели. Это позволит связать внутреннее устройство конвейера с тем, как его возможности выглядят и управляются в работающем приложении.

View the original on Хабр →

KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.