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

Запуск кількох потоків в одному графі

ПолеЗначення
КатегоріяГрафи та конвеєри
СкладністьПросунутий
Орієнтовний час читання20-25 minutes
Міткиgraph, multistream, scheduler, join

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

Кожен зразок, який ви передаєте, містить stream_id та frame_id. Політика об’єднання ByFrame чекає, поки обидва іменовані входи (left та right) не передадуть зразок з однаковим frame_id, після чого генерує рівно один об’єднаний пакет. В результаті ви створите граф об’єднання, розподілите детерміноване навантаження для кожного потоку/кадру через два входи та об’єднаєте пакети, перевіривши кількість вихідних даних і те, що кожен пакет містить два поля.

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

Створення графа об’єднання

graphs::Combine (C++) / graphs.combine (Python) повертає звичайний загальнодоступний фрагмент Graph — в ньому немає нічого особливого, окрім його форми: два іменовані вхідні потоки, один іменований вихідний потік і політика об’єднання. Ми передаємо ["left", "right"] як імена вхідних потоків, "combined" як ім’я вихідного потоку та CombinePolicy.ByFrame для вибору відповідності за ідентифікатором кадру. Друк describe() показує отриману топологію, а build() перетворює опис на виконуваний об’єкт. За замовчуванням граф працює асинхронно, тому кожен потік може незалежно просуватися вперед.

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

CombinePolicy.ByFrame здійснює відповідність за Sample.frame_id; CombinePolicy.ByPts є альтернативою, яка здійснює відповідність за часовими мітками презентації (Sample.pts_ns), коли кадри не мають чіткого індексу кадру.

tutorials/015_run_multiple_streams/run_multiple_streams.cpp
simaai::neat::Graph graph = simaai::neat::graphs::Combine({"left", "right"}, "combined",
simaai::neat::CombinePolicy::ByFrame);

std::cout << graph.describe() << "\n";

const int expected = streams * frames;
simaai::neat::Run run = graph.build();

Передайте потоки

Тепер ми обробляємо навантаження. Для кожного кадру та кожного потоку ми синтезуємо невеликий детермінований RGB-зразок, позначений його stream_id та унікальним frame_id, а потім передаємо його в обидва вказані вхідні потоки. Оскільки ідентифікатори обчислюються детерміновано (frame * streams + sid), з’єднання має однозначне відповідне поєднання для пошуку: кадр left N завжди має відповідний кадр right N. Після відповідного right передавання ми очищаємо об’єднаний вихід цієї пари, перш ніж перейти до наступної пари.

Кожен зразок явно створюється як Sample, що містить Tensor (HWC, UInt8, RGB) із встановленими frame_id та stream_id; run.push("left", sample) повертає булеве значення, яке слід перевірити на відповідність run.last_error().

tutorials/015_run_multiple_streams/run_multiple_streams.cpp
if (!run.push("left", make_rgb_sample(std::to_string(sid), logical_frame))) {
throw std::runtime_error("left push failed: " + run.last_error());
}
if (!run.push("right", make_rgb_sample(std::to_string(sid), logical_frame))) {
throw std::runtime_error("right push failed: " + run.last_error());
}

Витягуємо кожен об’єднаний пакет

Відразу після додавання кожної пари відповідних елементів, ми один раз отримуємо дані з вказаного вихідного потоку "combined". Кожне успішне отримання даних повертає пакет, який середовище виконання згенерувало після того, як обидва вхідні потоки надали відповідний кадр. Очищення під час генерації запобігає заповненню обмеженої вихідної черги та поширенню зворотного тиску на вхідну сторону. Обидва приклади підтверджують, що кожен пакет містить два об’єднані поля, а потім викликають close(), щоб коректно завершити виконання. Очікувана кількість пакетів дорівнює streams * frames, що підтверджує відсутність втрачених пар.

run.pull("combined", timeout_ms) повертає необов’язковий пакет; ми зчитуємо bundle.stream_id та bundle.fields.size() і перевіряємо, чи кожен пакет містить два поля.

tutorials/015_run_multiple_streams/run_multiple_streams.cpp
auto maybe_bundle = run.pull("combined", /*timeout_ms=*/2000);
if (!maybe_bundle.has_value()) {
throw std::runtime_error("timed out waiting for combined output: " + run.last_error());
}
const auto& bundle = *maybe_bundle;
const int fields = static_cast<int>(bundle.fields.size());
if (fields != 2)
throw std::runtime_error("joined bundle should contain two fields");
if (first_fields < 0)
first_fields = fields;
++received;
if (received <= 4) {
std::cout << "bundle stream=" << bundle.stream_id << " fields=" << fields << "\n";
}

Запуск

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

C++ (prebuilt):

./lib/sima-neat/tutorials/tutorial_015_run_multiple_streams \
--streams 8 --frames 4

C++ (build from source):

./build.sh --target tutorial_015_run_multiple_streams
./build/tutorials-standalone/tutorial_015_run_multiple_streams \
--streams 8 --frames 4

Очікуваний результат (збірка C++ також виводить опис графа; обидві збірки виводять перші кілька пакетів):

received=32 fields=2
[OK] 015_run_multiple_streams

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

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

Показати повні програми
tutorials/015_run_multiple_streams/run_multiple_streams.cpp
// Multistream public Graph: named inputs -> Combine(ByFrame) -> named output bundle.
//
// Usage:
// tutorial_015_run_multiple_streams [--streams 8] [--frames 4]

#include "neat.h"

#include <cstdint>
#include <iostream>
#include <stdexcept>
#include <string>
#include <utility>
#include <vector>

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

std::vector<int64_t> contiguous_strides_bytes(const std::vector<int64_t>& shape,
int64_t elem_bytes) {
std::vector<int64_t> strides(shape.size(), 0);
int64_t stride = elem_bytes;
for (int i = static_cast<int>(shape.size()) - 1; i >= 0; --i) {
strides[static_cast<size_t>(i)] = stride;
stride *= shape[static_cast<size_t>(i)];
}
return strides;
}

simaai::neat::Sample make_rgb_sample(const std::string& stream_id, int frame_id) {
const int w = 8;
const int h = 6;
const int c = 3;
const std::size_t bytes = static_cast<std::size_t>(w) * h * c;

simaai::neat::Tensor t;
t.device = {simaai::neat::DeviceType::CPU, 0};
t.dtype = simaai::neat::TensorDType::UInt8;
t.layout = simaai::neat::TensorLayout::HWC;
t.shape = {h, w, c};
t.semantic.image = simaai::neat::ImageSpec{simaai::neat::ImageSpec::PixelFormat::RGB, ""};
t.storage = simaai::neat::make_cpu_owned_storage(bytes);
t.strides_bytes = contiguous_strides_bytes(t.shape, 1);
t.read_only = false;
{
auto map = t.map(simaai::neat::MapMode::Write);
auto* p = static_cast<std::uint8_t*>(map.data);
for (std::size_t i = 0; i < bytes; ++i)
p[i] = static_cast<std::uint8_t>(i % 255);
}
t.read_only = true;

simaai::neat::Sample sample;
sample.kind = simaai::neat::SampleKind::Tensor;
sample.tensor = std::move(t);
sample.frame_id = frame_id;
sample.stream_id = stream_id;
return sample;
}

} // namespace

int main(int argc, char** argv) {
try {
const int streams = parse_int_arg(argc, argv, "--streams", 8);
const int frames = parse_int_arg(argc, argv, "--frames", 4);

// CORE LOGIC
// `graphs::Combine` is a normal public Graph fragment. It declares two
// named inputs ("left", "right") and one named output ("combined"). ByFrame
// means the runtime emits one bundle only after both inputs have delivered
// samples with the same Sample::frame_id.
simaai::neat::Graph graph = simaai::neat::graphs::Combine({"left", "right"}, "combined",
simaai::neat::CombinePolicy::ByFrame);

std::cout << graph.describe() << "\n";

const int expected = streams * frames;
simaai::neat::Run run = graph.build();

int received = 0;
int first_fields = -1;
for (int frame = 0; frame < frames; ++frame) {
for (int sid = 0; sid < streams; ++sid) {
const int logical_frame = frame * streams + sid;

if (!run.push("left", make_rgb_sample(std::to_string(sid), logical_frame))) {
throw std::runtime_error("left push failed: " + run.last_error());
}
if (!run.push("right", make_rgb_sample(std::to_string(sid), logical_frame))) {
throw std::runtime_error("right push failed: " + run.last_error());
}

auto maybe_bundle = run.pull("combined", /*timeout_ms=*/2000);
if (!maybe_bundle.has_value()) {
throw std::runtime_error("timed out waiting for combined output: " + run.last_error());
}
const auto& bundle = *maybe_bundle;
const int fields = static_cast<int>(bundle.fields.size());
if (fields != 2)
throw std::runtime_error("joined bundle should contain two fields");
if (first_fields < 0)
first_fields = fields;
++received;
if (received <= 4) {
std::cout << "bundle stream=" << bundle.stream_id << " fields=" << fields << "\n";
}
}
}

run.close();

if (received != expected)
throw std::runtime_error("expected=" + std::to_string(expected) +
" received=" + std::to_string(received));
if (first_fields != 2)
throw std::runtime_error("join should emit a two-field bundle");

std::cout << "received=" << received << " fields=" << first_fields << "\n";
std::cout << "[OK] 015_run_multiple_streams\n";
return 0;
} catch (const std::exception& e) {
std::cerr << "[FAIL] " << e.what() << "\n";
return 1;
}
}

Джерело