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

Налаштування пропускної здатності та глибини черги

ПолеЗначення
КатегоріяГрафи та конвеєри
СкладністьПросунутий
Орієнтовний час читання15-20 minutes
Міткиperformance, tuning, async, queues

Налаштування продуктивності дає позитивний результат лише тоді, коли ваша базова модель коректності є стабільною; у цій главі ми припускаємо, що це так, і переходимо до налаштувань, які визначають, як працює асинхронний конвеєр, коли обсяг роботи перевищує швидкість його обробки. Ви встановите глибину черги, визначите, що відбуватиметься, коли черга заповниться, передасте детерміновану серію кадрів у неблокуючому режимі, обробите результати та переглянете звіт про вимірювання, який покаже, чи були втрачені дані та скільки часу знадобилося для обробки кожного кадру.

У результаті ви отримаєте робочу систему для вимірювання асинхронного виконання в умовах обмеженої пропускної здатності: кількість елементів у черзі, кількість втрачених елементів, кількість отриманих результатів, середню затримку та вартість передачі. Той самий цикл є основою для налаштування реального конвеєра відповідно до евристик, описаних у На практиці.

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

Налаштуйте параметри виконання

RunOptions визначає, як реалізується асинхронна поведінка під навантаженням. Ми встановлюємо queue_depth (кількість зразків, які може прийняти середовище виконання), overflow_policy (що відбувається, коли черга заповнюється — Block, KeepLatest або DropIncoming), output_memory = Owned (повернені тензори зберігають свої дані, тому вони залишаються доступними після отримання). Потім ми build() граф у Async режимі, що забезпечує виконання з незалежними сторонами виробника та споживача.

Політику обробки переповнення аналізують з --drop у simaai::neat::OverflowPolicy::{Block,KeepLatest,DropIncoming}; graph.build(input, opt) повертає дескриптор виконання.

tutorials/016_tune_throughput_and_queues/tune_throughput_and_queues.cpp
simaai::neat::RunOptions opt;
opt.queue_depth = queue_depth;
opt.overflow_policy = parse_drop_policy(argc, argv);
opt.output_memory = simaai::neat::OutputMemory::Owned;

auto run = graph.build(std::vector<cv::Mat>{rgb}, opt);

Передайте робоче навантаження та звільніть

Тут реалізується політика черги. Ми викликаємо try_push(...) у тісному циклі — неблокуюче додавання, яке просто повертає, чи було прийнято зразок, тому повна черга за умови DropIncoming/ KeepLatest відображається як відхилені спроби додавання, а не як зупинка. Після сплеску ми викликаємо close_input(), щоб повідомити про відсутність подальших вхідних даних, а потім очищаємо сторону споживача за допомогою циклу pull(...), поки він не поверне порожній результат. Поєднання try_push з close_input разом із циклом очищення є стандартним неблокуючим асинхронним шаблоном.

tutorials/016_tune_throughput_and_queues/tune_throughput_and_queues.cpp
// try_push never blocks; pair it with close_input + drain pull loop.
simaai::neat::MeasureOptions measure_opt;
measure_opt.title = "tutorial 016 throughput";
auto scope = run.start_measurement(measure_opt);
for (int i = 0; i < iters; ++i)
(void)run.try_push(std::vector<cv::Mat>{rgb});
run.close_input();

int pulled = 0;
while (run.pull(/*timeout_ms=*/1000).has_value())
++pulled;
const auto measured = scope.stop();
if (measured.counters.inputs_enqueued <= 0 || pulled <= 0)
throw std::runtime_error("throughput run produced no measured inputs/outputs");

Перегляньте звіт про вимірювання

Після завершення виконання ми зупиняємо область вимірювання. Група counters у звіті надає дані зі сторони середовища виконання: кількість доданих вхідних даних, кількість відхилених вхідних даних, кількість отриманих вихідних даних, тоді як input надає дані зі сторони додавання, такі як середня вартість додавання та кількість повторних переговорів щодо вхідних даних. Разом вони показують, чи працювала глибина вашої черги та політика обробки переповнення так, як ви очікували: чи відбувалися втрати кадрів, чи збільшувалася затримка, чи була вартість шляху додавання низькою.

tutorials/016_tune_throughput_and_queues/tune_throughput_and_queues.cpp
std::cout << "inputs_enqueued=" << measured.counters.inputs_enqueued << "\n";
std::cout << "inputs_dropped=" << measured.counters.inputs_dropped << "\n";
std::cout << "outputs_pulled=" << pulled << "\n";
std::cout << "avg_latency_ms=" << measured.end_to_end.avg_ms << "\n";
std::cout << "avg_push_us=" << measured.input.avg_push_us << "\n";
std::cout << "renegotiations=" << measured.input.renegotiations << "\n";

Запуск

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

C++ (prebuilt):

./lib/sima-neat/tutorials/tutorial_016_tune_throughput_and_queues \
--iters 32 --queue 4 --drop block

C++ (build from source):

./build.sh --target tutorial_016_tune_throughput_and_queues
./build/tutorials-standalone/tutorial_016_tune_throughput_and_queues \
--iters 32 --queue 4 --drop block

Очікуваний результат (точні значення та час виконання залежать від хоста та політики):

inputs_enqueued=32
inputs_dropped=0
outputs_pulled=32
avg_latency_ms=0.42
avg_push_us=18.0
renegotiations=0
[OK] 016_tune_throughput_and_queues

(Під час збірки Python виводяться ті самі ключі, але без останнього рядка [OK]).

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

На практиці

Практичні поради щодо визначення розміру черги, політик відхилення, попередньо встановлених параметрів і забезпечення безпеки терміну дії вихідних даних.

Визначення розміру черги (queue_depth)

Евристичні правила:

  • Почніть з queue_depth = 4–16 для конвеєрів із низькою затримкою.
  • Збільште розмір черги, якщо ваш генератор даних працює з переривчастим потоком або якщо наступні елементи мають змінну затримку (декодування/MLA/постобробка).
  • Зберігайте невеликий розмір черги, якщо вам потрібні найсвіжіші кадри (наприклад, для попереднього перегляду відео в реальному часі).

Політика обробки переповнення (RunOptions::overflow_policy)

  • Block: найбезпечніший варіант для забезпечення коректності; генератор даних чекає, коли черга звільниться.
  • DropIncoming: зберігає дані в черзі, відкидає вхідні зразки, коли черга заповнена.
  • KeepLatest: надає перевагу найсвіжішим кадрам, відкидає найстаріші зразки в черзі.

Для потокового відео в реальному часі, KeepLatest зазвичай забезпечує найменшу загальну затримку.

Попередньо встановлені параметри та повторне узгодження

Використовуйте RunOptions::preset для керування компромісами між затримкою та безпекою:

  • Realtime: найменша затримка, агресивна поведінка щодо забезпечення свіжості даних.
  • Balanced: починає використовувати копіювання без буферизації, коли це можливо, виконує перевірки під час запуску та переходить у режим копіювання, якщо виникають проблеми з надійністю.
  • Reliable: консервативна поведінка та стабільне володіння вихідними даними.

Перегляд форми вхідних даних відбувається автоматично для динамічних вхідних даних (лічильник renegotiations, наведений вище, показує, як часто це відбувалося).

Час життя вихідних даних (output_memory)

  • output_memory = Owned: повернутий Tensor володіє своїми даними.
  • output_memory = ZeroCopy: тензор може посилатися на буфери середовища виконання, які повторно використовуються після отримання даних.
  • output_memory = Auto: середовище виконання спочатку намагається використовувати механізм без копіювання, а якщо це неможливо, переходить до механізму, коли дані належать тензору, якщо це необхідно для забезпечення надійності.

Якщо вам потрібно зберегти дані тензора після поточного кроку, викличте clone() або cpu().contiguous().

Безпека пулу буферів

  • RunAdvancedOptions::max_input_bytes встановлює жорстку верхню межу для виділення вхідних буферів.
  • Якщо потрібен більший буфер, середовище виконання швидко завершує роботу з явним повідомленням про помилку.

Використовуйте ці параметри, щоб захистити довготривалі процеси від необмеженого виділення пам’яті, коли розмір вхідних даних змінюється.

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

Показати повні програми
tutorials/016_tune_throughput_and_queues/tune_throughput_and_queues.cpp
// Tune async Graph throughput via RunOptions and MeasureReport.
//
// Usage:
// tutorial_016_tune_throughput_and_queues [--iters 32] [--queue 4] [--drop block|latest|incoming]

#include "neat.h"

#include <opencv2/core.hpp>

#include <iostream>
#include <stdexcept>
#include <string>

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);
}

simaai::neat::OverflowPolicy parse_drop_policy(int argc, char** argv) {
std::string mode;
if (!get_arg(argc, argv, "--drop", mode))
return simaai::neat::OverflowPolicy::Block;
if (mode == "latest")
return simaai::neat::OverflowPolicy::KeepLatest;
if (mode == "incoming")
return simaai::neat::OverflowPolicy::DropIncoming;
return simaai::neat::OverflowPolicy::Block;
}

} // namespace

int main(int argc, char** argv) {
try {
const int iters = parse_int_arg(argc, argv, "--iters", 32);
const int queue_depth = parse_int_arg(argc, argv, "--queue", 4);

cv::Mat rgb(120, 160, CV_8UC3, cv::Scalar(70, 20, 200));
if (!rgb.isContinuous())
rgb = rgb.clone();

simaai::neat::Graph graph;
simaai::neat::InputOptions in;
in.format = "RGB";
in.width = rgb.cols;
in.height = rgb.rows;
in.depth = rgb.channels();
in.is_live = true;
graph.add(simaai::neat::nodes::Input(in));
graph.add(simaai::neat::nodes::Output());

// CORE LOGIC
// RunOptions controls how the async runner buffers and drops frames.
simaai::neat::RunOptions opt;
opt.queue_depth = queue_depth;
opt.overflow_policy = parse_drop_policy(argc, argv);
opt.output_memory = simaai::neat::OutputMemory::Owned;

auto run = graph.build(std::vector<cv::Mat>{rgb}, opt);

// try_push never blocks; pair it with close_input + drain pull loop.
simaai::neat::MeasureOptions measure_opt;
measure_opt.title = "tutorial 016 throughput";
auto scope = run.start_measurement(measure_opt);
for (int i = 0; i < iters; ++i)
(void)run.try_push(std::vector<cv::Mat>{rgb});
run.close_input();

int pulled = 0;
while (run.pull(/*timeout_ms=*/1000).has_value())
++pulled;
const auto measured = scope.stop();
if (measured.counters.inputs_enqueued <= 0 || pulled <= 0)
throw std::runtime_error("throughput run produced no measured inputs/outputs");

std::cout << "inputs_enqueued=" << measured.counters.inputs_enqueued << "\n";
std::cout << "inputs_dropped=" << measured.counters.inputs_dropped << "\n";
std::cout << "outputs_pulled=" << pulled << "\n";
std::cout << "avg_latency_ms=" << measured.end_to_end.avg_ms << "\n";
std::cout << "avg_push_us=" << measured.input.avg_push_us << "\n";
std::cout << "renegotiations=" << measured.input.renegotiations << "\n";
std::cout << "[OK] 016_tune_throughput_and_queues\n";
return 0;
} catch (const std::exception& e) {
std::cerr << "[FAIL] " << e.what() << "\n";
return 1;
}
}

Джерело