Перейти до основного вмісту

Запуск висновків асинхронно

ПолеЗначення
КатегоріяМоделі та інференс
СкладністьПочатковий
Орієнтовний час читання10-15 minutes
Міткиasync, push-pull, throughput, runtime

Розділ 001 запустив модель з одним синхронним викликом: передає один кадр, блокує, поки не отримає результат. Це просто, але неефективно — потік, який генерує вхідні дані, і потік, який обробляє вихідні дані, є одним і тим же потоком, тому вони ніколи не можуть працювати паралельно. У цьому розділі використовується точно така ж модель ResNet-50, але її перетворюють на орієнтований на пропускну здатність конвеєр, розділивши ці дві задачі.

Механізм полягає в асинхронному Run: ви build() модель у Graph в асинхронному режимі (Async), а потім керуєте нею за допомогою двох незалежних викликів — push(...) від генератора та pull(...) від споживача. В результаті ви отримаєте потік-генератор, який передає кадри якомога швидше, з урахуванням можливостей середовища виконання, а головний потік витягує прогнози, і остаточний рядок pushed=N pulled=N, який доводить, що нічого не було втрачено.

Покроковий огляд

Завантаження моделі

Ми починаємо точно так само, як у розділі 001 — створюємо Model з архіву — але тут ми також визначаємо RouteOptions зі встановленими параметрами include_input та include_output. Ці прапорці вказують моделі, що вона повинна надавати власні вхідні та вихідні межі, коли її включають у граф, щоб навколишній конвеєр міг передавати кадри та витягувати тензори.

tutorials/002_run_inference_async/run_inference_async.cpp
simaai::neat::Model model(model_path, build_options(size));
simaai::neat::Model::RouteOptions route_opt;
route_opt.include_input = true;
route_opt.include_output = true;

Створення асинхронного конвеєра

Model не можна безпосередньо використовувати з push/pull; для цього потрібен Run. Ми обертаємо модель у новий Graph за допомогою graph.add(model.graph(route_opt)), а потім build(...) його з використанням репрезентативного кадру. Передача кадру-зразка дозволяє build() заздалегідь узгодити конкретні форми тензорів. Повернений Run є об’єктом, яким користуватимуться обидва потоки.

tutorials/002_run_inference_async/run_inference_async.cpp
simaai::neat::Graph graph;
graph.add(model.graph(route_opt));

auto run = graph.build(std::vector<cv::Mat>{frames.front()});

Передача кадрів від генератора

Єдина задача генератора — передавати вхідні дані. Ми запускаємо потік, який перебирає підготовлені кадри, викликає push(...) для кожного, а потім викликає close_input(), щоб повідомити, що більше кадрів не буде — цей сигнал дає споживачу знати, коли зупинитися. Оскільки генератор працює незалежно, він не чекає на жоден результат, перш ніж надсилати наступний кадр.

std::thread запускає цикл; атомний лічильник pushed і прапорець producer_done оновлюються в процесі, щоб головний потік міг спостерігати за прогресом без блокування.

tutorials/002_run_inference_async/run_inference_async.cpp
std::atomic<int> pushed{0};
std::atomic<bool> producer_done{false};
std::thread producer([&]() {
for (const cv::Mat& f : frames) {
run.push(std::vector<cv::Mat>{f});
pushed.fetch_add(1, std::memory_order_relaxed);
}
run.close_input();
producer_done.store(true);
});

Отримання результатів споживачем

Основний потік споживає дані. Він виконує цикл, викликаючи pull(timeout_ms=2000), що повертає наступний доступний результат або нічого, якщо протягом заданого часу очікування нічого не надійшло. Якщо отримано порожній результат, ми перевіряємо, чи завершив роботу виробник — якщо так, ми зупиняємося, інакше ми продовжуємо чекати. Кожен отриманий результат зводиться до індексу класу з найбільшим значенням (top-1) і виводиться на екран. Після завершення циклу ми завершуємо роботу виробника та підтверджуємо, що дані були pushed == pulled.

pull() повертає optional<Sample>; отримайте тензори за допомогою tensors_from_sample(...) перед читанням байтів.

tutorials/002_run_inference_async/run_inference_async.cpp
int pulled = 0;
while (pulled < n) {
auto out = run.pull(/*timeout_ms=*/2000);
if (!out.has_value()) {
if (producer_done.load())
break;
continue;
}
std::cout << "top1=" << top1_from_output(*out) << "\n";
++pulled;
}
producer.join();

Запуск

Запустіть програму, і ви повинні побачити один рядок top1= на кадр, за яким слідує підрахунок кількості відправлених/отриманих даних. Запустіть команди Python і C++ (попередньо скомпільовані) з кореневої Neat теки (директорії, що містить share/ і lib/); запустіть команди збірки з вихідного коду з кореневої теки репозиторію.

C++ (prebuilt):

./lib/sima-neat/tutorials/tutorial_002_run_inference_async \
--model /tmp/resnet_50.tar.gz --n 4

C++ (build from source):

./build.sh --target tutorial_002_run_inference_async
./build/tutorials-standalone/tutorial_002_run_inference_async \
--model /tmp/resnet_50.tar.gz --n 4

Очікуваний результат (точні індекси залежать від зображення; збірка C++ додає поле pushed=..., збірка Python виводить лише pulled=...):

top1=285
top1=285
top1=285
top1=285
pushed=4 pulled=4
[OK] 002_run_inference_async

Щоб інтегрувати вихідний код цього розділу у власний проєкт із використанням спеціального CMakeLists.txt (додаткова тека не потрібна), див. розділ Як запускати навчальні матеріали на головній сторінці.

На практиці

У цьому розділі використовується асинхронний механізм обміну даними. Щоб виміряти ту саму модель із детермінованими синтетичними вхідними даними, перейдіть до розділу Оцінка вашої моделі. Для повної моделі, що порівнює збірку та запуск, а також синхронний і асинхронний режими, а також повний набір RunOptions, див. розділ Створення вашого першого графа. Для глибини черги, політики обробки переповнення та вимірювання під навантаженням див. розділ Налаштування пропускної здатності та глибини черги.

Повний початковий код

Показати повні програми
tutorials/002_run_inference_async/run_inference_async.cpp
// Async push/pull: producer thread pushes frames, main thread pulls outputs.
//
// Usage:
// tutorial_002_run_inference_async --model /path/to/resnet_50.tar.gz [--image /path/to.jpg] [--n
// 4]

#include "neat.h"

#include <opencv2/core.hpp>
#include <opencv2/imgcodecs.hpp>
#include <opencv2/imgproc.hpp>

#include <atomic>
#include <cstring>
#include <exception>
#include <filesystem>
#include <iostream>
#include <stdexcept>
#include <string>
#include <thread>
#include <vector>

namespace fs = std::filesystem;

namespace {

bool get_arg(int argc, char** argv, const std::string& key, std::string& out) {
for (int i = 1; i + 1 < argc; ++i) {
if (key == argv[i]) {
out = argv[i + 1];
return true;
}
}
return false;
}

int parse_int_arg(int argc, char** argv, const std::string& key, int def) {
std::string value;
if (!get_arg(argc, argv, key, value))
return def;
return std::stoi(value);
}

cv::Mat load_rgb(const fs::path& image_path, int size) {
cv::Mat bgr = cv::imread(image_path.string(), cv::IMREAD_COLOR);
if (bgr.empty())
throw std::runtime_error("failed to read image: " + image_path.string());
if (bgr.cols != size || bgr.rows != size) {
cv::resize(bgr, bgr, cv::Size(size, size), 0, 0, cv::INTER_AREA);
}
cv::Mat rgb;
cv::cvtColor(bgr, rgb, cv::COLOR_BGR2RGB);
if (!rgb.isContinuous())
rgb = rgb.clone();
return rgb;
}

simaai::neat::Model::Options build_options(int size) {
simaai::neat::Model::Options opt;
opt.preprocess.color_convert.input_format = simaai::neat::PreprocessColorFormat::RGB;
opt.preprocess.input_max_width = size;
opt.preprocess.input_max_height = size;
opt.preprocess.input_max_depth = 3;
opt.preprocess.normalize.mean = {0.485f, 0.456f, 0.406f};
opt.preprocess.normalize.stddev = {0.229f, 0.224f, 0.225f};
return opt;
}

int top1_from_output(const simaai::neat::Sample& out) {
if (simaai::neat::tensors_from_sample(out, true).empty())
throw std::runtime_error("no tensor output");
const simaai::neat::Mapping m = simaai::neat::tensors_from_sample(out, true).front().map_read();
const size_t n = m.size_bytes / sizeof(float);
const float* p = reinterpret_cast<const float*>(m.data);
int best = 0;
for (size_t i = 1; i < n && i < 1000; ++i) {
if (p[i] > p[best])
best = static_cast<int>(i);
}
return best;
}

} // namespace

int main(int argc, char** argv) {
try {
std::string model_path, image;
if (!get_arg(argc, argv, "--model", model_path)) {
std::cerr
<< "Usage: tutorial_002_run_inference_async --model <path> [--image <path>] [--n <n>]\n";
return 1;
}
get_arg(argc, argv, "--image", image);
const int n = parse_int_arg(argc, argv, "--n", 4);
const int size = 224;

cv::Mat frame = image.empty() ? cv::Mat(size, size, CV_8UC3, cv::Scalar(99, 99, 99))
: load_rgb(image, size);
std::vector<cv::Mat> frames(n, frame);

// CORE LOGIC
// Build a Graph around the model and run it async: one producer thread pushes,
// the main thread pulls outputs.
simaai::neat::Model model(model_path, build_options(size));
simaai::neat::Model::RouteOptions route_opt;
route_opt.include_input = true;
route_opt.include_output = true;

simaai::neat::Graph graph;
graph.add(model.graph(route_opt));

auto run = graph.build(std::vector<cv::Mat>{frames.front()});

std::atomic<int> pushed{0};
std::atomic<bool> producer_done{false};
std::thread producer([&]() {
for (const cv::Mat& f : frames) {
run.push(std::vector<cv::Mat>{f});
pushed.fetch_add(1, std::memory_order_relaxed);
}
run.close_input();
producer_done.store(true);
});

int pulled = 0;
while (pulled < n) {
auto out = run.pull(/*timeout_ms=*/2000);
if (!out.has_value()) {
if (producer_done.load())
break;
continue;
}
std::cout << "top1=" << top1_from_output(*out) << "\n";
++pulled;
}
producer.join();

std::cout << "pushed=" << pushed.load() << " pulled=" << pulled << "\n";
if (pulled != n)
throw std::runtime_error("pulled=" + std::to_string(pulled) +
" != pushed=" + std::to_string(pushed.load()));
std::cout << "[OK] 002_run_inference_async\n";
return 0;
} catch (const std::exception& e) {
std::cerr << "[FAIL] " << e.what() << "\n";
return 1;
}
}

Джерело