調整輸送量和佇列深度
| 欄位 | 值 |
|---|---|
| 類別 | 圖形與管線 |
| 難度 | 進階 |
| 預估閱讀時間 | 15-20 minutes |
| 標籤 | performance, tuning, async, queues |
效能調整只有在您的正確性基準穩定時才會有所幫助;本章假設已達到穩定狀態,並著重於決定非同步管線在工作量以超過其處理能力的速度到達時的行為方式。您將設定佇列深度,選擇在該佇列已滿時發生的情況,以非阻塞方式推送一組確定的影格,排出結果,並讀取測量報告,以了解您是否丟棄了任何內容以及每個影格需要多長時間。
到最後,您將擁有一個用於測量在反壓情況下非同步運行的可用框架:佇列計數、丟棄計數、提取的輸出、平均延遲和推送成本。相同的迴圈是針對真實管線,根據 實際應用 中的啟發式方法進行調整的基礎。
操作指南
設定執行選項
RunOptions 是決定負載下非同步行為的地方。我們設定 queue_depth(執行階段接受的同時進行的樣本數量)、overflow_policy(當該 佇列已滿時發生的情況——Block、KeepLatest 或 DropIncoming)、output_memory = Owned(傳回的張量擁有其資料,因此在提取後仍會保留)。然後,我們以 Async 模式 build() 圖,這為我們提供了一個具有獨立的生產者和消費者端的運行。
C++: 超出限制策略從 --drop 解析為 simaai::neat::OverflowPolicy::{Block,KeepLatest,DropIncoming};graph.build(input, opt) 傳回運行句柄。
Python: 策略使用 getattr(pyneat.OverflowPolicy, ...) 解析;graph.build([tensor], opt) 傳回運行句柄。
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 結合,再加上一個排出迴圈,是標準的非阻塞非同步模式。
// 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 提供推送端的指標,例如平均推送成本和重新協商的輸入。這些指標可以幫助您判斷,您的佇列深度和溢出策略是否達到預期效果:是否發生封包遺失、延遲是否增加,以及推送路徑的成本是否降低。
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設定輸入緩衝區分配的硬性上限。- 如果需要更大的緩衝區,則執行階段會立即失敗,並顯示明確的錯誤。
當輸入大小發生變化時,使用這些設定來保護長時間執行的程序,使其免受無限分配的影響。
完整原始碼
顯示完整原始碼程式
// 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;
}
}