Workflowsベンチマーク
直接のモデル推論と、同じモデルをWorkflowでラップした場合のレイテンシを比較します。
最終更新
役に立ちましたか?
役に立ちましたか?
import os
import statistics
import time
import argparse
import csv
import supervision as sv
from inference import get_model
from inference.core.env import WORKFLOWS_MAX_CONCURRENT_STEPS, MAX_ACTIVE_MODELS
from inference.core.managers.base import ModelManager
from inference.core.managers.decorators.fixed_size_cache import WithFixedSizeCache
from inference.core.registries.roboflow import RoboflowModelRegistry
from inference.core.workflows.core_steps.common.entities import StepExecutionMode
from inference.core.workflows.execution_engine.core import ExecutionEngine
from inference.models.utils import ROBOFLOW_MODEL_TYPES
def build_workflow(model_id: str) -> dict:
"""指定されたモデル(検出または分類)用の最小限のワークフロー定義を構築します。"""
if "classifiers" in model_id or "classification" in model_id:
step_type = "RoboflowClassificationModel"
else:
step_type = "RoboflowObjectDetectionModel"
return {
"version": "1.0",
"inputs": [
{"type": "WorkflowImage", "name": "image"},
],
"steps": [
{
"type": step_type,
"name": "model_step",
"image": "$inputs.image",
"model_id": model_id,
}
],
"outputs": [
{
"type": "JsonField",
"name": "predictions",
"selector": "$steps.model_step.predictions",
},
],
}
def main():
parser = argparse.ArgumentParser(description="直接推論とワークフロー推論のレイテンシーをベンチマークします。")
parser.add_argument("--iterations", type=int, default=10, help="各メソッドごとの計測反復回数(既定: 10)")
args = parser.parse_args()
# モデル、benchmark.py と同じリスト
models = [
"classifiers/3",
"yolo26n-640", "yolo26s-640", "yolo26m-640", "yolo26l-640", "yolo26x-640",
"rfdetr-nano", "rfdetr-small", "rfdetr-medium", "rfdetr-large", "rfdetr-xlarge", "rfdetr-2xlarge",
]
# テスト画像を 1 回だけダウンロード
print("テスト画像をダウンロード中...")
url = "https://media.roboflow.com/inference/people-walking.jpg"
image_np = sv.load_image_from_url(url)
print("画像の準備ができました。")
# Workflow エンジン用の共有モデルマネージャーを初期化(モデル間で再利用)
model_registry = RoboflowModelRegistry(ROBOFLOW_MODEL_TYPES)
model_manager = ModelManager(model_registry=model_registry)
model_manager = WithFixedSizeCache(model_manager, max_size=MAX_ACTIVE_MODELS)
workflow_init_parameters = {
"workflows_core.model_manager": model_manager,
"workflows_core.step_execution_mode": StepExecutionMode.LOCAL,
}
results_file = os.path.join(os.path.dirname(os.path.abspath(__file__)), "benchmark_inference_vs_workflows.csv")
fieldnames = [
"model_id",
"avg_latency_direct_ms", "min_direct_ms", "max_direct_ms", "stddev_direct_ms",
"avg_latency_workflow_ms", "min_workflow_ms", "max_workflow_ms", "stddev_workflow_ms",
]
with open(results_file, "w", newline="") as f:
writer = csv.DictWriter(f, fieldnames=fieldnames)
writer.writeheader()
print(f"\n各メソッドにつき {args.iterations} 回の反復を {len(models)} 個のモデルで実行します...\n")
for model_id in models:
print(f"--- モデル: {model_id} ---")
# ── 1. 直接推論 ──────────────────────────────────────────────
direct_latencies = None
try:
model = get_model(model_id=model_id)
# ウォームアップ
model.infer(image_np)
direct_latencies = []
for i in range(args.iterations):
t0 = time.perf_counter()
model.infer(image_np)
t1 = time.perf_counter()
ms = (t1 - t0) * 1000.0
direct_latencies.append(ms)
print(f" 直接 [{i+1:>2}/{args.iterations}]: {ms:.2f} ms")
avg = statistics.mean(direct_latencies)
mn = min(direct_latencies)
mx = max(direct_latencies)
sd = statistics.stdev(direct_latencies) if len(direct_latencies) > 1 else 0.0
print(f" 直接 → 平均={avg:.2f} 最小={mn:.2f} 最大={mx:.2f} 標準偏差={sd:.2f} ms")
except Exception as e:
print(f" 直接失敗: {e}")
# ── 2. ワークフロー推論 ────────────────────────────────────────────
workflow_latencies = None
try:
workflow_def = build_workflow(model_id)
engine = ExecutionEngine.init(
workflow_definition=workflow_def,
init_parameters=workflow_init_parameters,
max_concurrent_steps=WORKFLOWS_MAX_CONCURRENT_STEPS,
)
# ウォームアップ
engine.run(runtime_parameters={"image": [image_np]})
workflow_latencies = []
for i in range(args.iterations):
t0 = time.perf_counter()
engine.run(runtime_parameters={"image": [image_np]})
t1 = time.perf_counter()
ms = (t1 - t0) * 1000.0
workflow_latencies.append(ms)
print(f" ワークフロー [{i+1:>2}/{args.iterations}]: {ms:.2f} ms")
avg = statistics.mean(workflow_latencies)
mn = min(workflow_latencies)
mx = max(workflow_latencies)
sd = statistics.stdev(workflow_latencies) if len(workflow_latencies) > 1 else 0.0
print(f" ワークフロー → 平均={avg:.2f} 最小={mn:.2f} 最大={mx:.2f} 標準偏差={sd:.2f} ms")
except Exception as e:
print(f" ワークフロー失敗: {e}")
# ── 行を書き込み ────────────────────────────────────────────────────────
def fmt(vals, fn):
return round(fn(vals), 2) if vals else "失敗"
row = {
"model_id": model_id,
"avg_latency_direct_ms": fmt(direct_latencies, statistics.mean),
"min_direct_ms": fmt(direct_latencies, min),
"max_direct_ms": fmt(direct_latencies, max),
"stddev_direct_ms": fmt(direct_latencies, lambda v: statistics.stdev(v) if len(v) > 1 else 0.0),
"avg_latency_workflow_ms": fmt(workflow_latencies, statistics.mean),
"min_workflow_ms": fmt(workflow_latencies, min),
"max_workflow_ms": fmt(workflow_latencies, max),
"stddev_workflow_ms": fmt(workflow_latencies, lambda v: statistics.stdev(v) if len(v) > 1 else 0.0),
}
with open(results_file, "a", newline="") as f:
writer = csv.DictWriter(f, fieldnames=fieldnames)
writer.writerow(row)
print(f" 保存しました。\n")
print(f"完了しました!結果: {results_file}")
if __name__ == "__main__":
main()