import os
import sys
import time
import multiprocessing
import argparse
import json
try:
# 导入华为 MindX SDK (mxVision)
from StreamManagerApi import StreamManagerApi, MxDataInput
except ImportError:
print("[Error] 无法导入 StreamManagerApi 或 MxDataInput。")
print("要在 Python 中测试 VDEC 硬件解码,必须正确安装华为的 MindX SDK (mxVision)。")
print("请确认服务器的环境变量是否已加载 (source set_env.sh),或检查 libffi 动态库冲突。")
sys.exit(1)
def build_pipeline_config(device_id, model_path, stream_id):
pipeline_name = f"vdec_infer_pipeline_{stream_id}"
pipeline = {
pipeline_name: {
"stream_config": {
"deviceId": str(device_id)
},
"appsrc0": {
"factory": "appsrc",
"props": {
"blocksize": "409600"
},
"next": "h264parse0"
},
"h264parse0": {
"factory": "h264parse",
"next": "mxpi_videodecoder0"
},
# mxpi_videodecoder0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0113.html
"mxpi_videodecoder0": {
"factory": "mxpi_videodecoder",
"props": {},
"next": "mxpi_imageresize0"
},
# mxpi_imageresize0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0111.html
"mxpi_imageresize0": {
"factory": "mxpi_imageresize",
"props": {
"resizeType": "Resizer_Stretch",
"resizeHeight": "640",
"resizeWidth": "640"
},
"next": "mxpi_tensorinfer0"
},
# mxpi_tensorinfer0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0122.html
"mxpi_tensorinfer0": {
"factory": "mxpi_tensorinfer",
"props": {
"dataSource": "mxpi_imageresize0",
"modelPath": model_path
},
"next": "appsink0"
},
"appsink0": {
"factory": "appsink"
}
}
}
return json.dumps(pipeline).encode('utf-8'), pipeline_name.encode('utf-8')
def run_single_stream(stream_id, video_path, model_path, device_id, test_duration=60):
"""
单路视频流的压测逻辑(运行在独立子进程中)。
参数:
stream_id : 当前流的编号,用于日志区分和唯一 stream_name 生成
video_path : H.264 裸流文件的绝对路径(由主进程传入)
model_path : OM 模型文件的绝对路径(由主进程传入)
device_id : NPU 设备 ID
test_duration : 压测持续时间(秒)
"""
prefix = f"[Stream {stream_id}]"
# 1. 初始化 StreamManager
stream_manager = StreamManagerApi()
ret = stream_manager.InitManager()
if ret != 0:
print(f"{prefix} 初始化 StreamManager 失败, ret={ret}")
return
# 2. 构建并创建 Pipeline 流
pipeline_conf, stream_name = build_pipeline_config(device_id, model_path, stream_id)
ret = stream_manager.CreateMultipleStreams(pipeline_conf)
if ret != 0:
print(f"{prefix} 创建 Pipeline 失败, ret={ret}")
print(f"{prefix} 请确认模型路径是否正确: {model_path}")
stream_manager.DestroyAllStreams()
return
in_plugin_id = 0
print(f"{prefix} Pipeline 创建成功,开始压测... (stream_name={stream_name.decode()})")
# 3. 读取本地 H.264 裸流文件
# 注意:必须是用 ffmpeg 抽取的 .h264 裸流,不能直接使用 .mp4 容器格式
# 抽取命令示例:ffmpeg -i input.mp4 -codec copy -bsf:v h264_mp4toannexb output.h264
try:
with open(video_path, "rb") as f:
video_bytes = f.read()
except Exception as e:
print(f"{prefix} 视频文件读取失败: {e}")
stream_manager.DestroyAllStreams()
return
if len(video_bytes) == 0:
print(f"{prefix} 视频文件为空,退出。")
stream_manager.DestroyAllStreams()
return
start_time = time.time()
frame_count = 0
error_count = 0
chunk_size = 409600 # 必须与 pipeline 中 appsrc0 的 blocksize 匹配 (400KB)
total_len = len(video_bytes)
offset = 0
# 4. 循环向 NPU 灌入数据并拉取推理结果
while time.time() - start_time < test_duration:
# 循环读取视频字节(到末尾后从头开始,模拟持续流)
end_offset = offset + chunk_size
if end_offset > total_len:
data = video_bytes[offset:] + video_bytes[:end_offset - total_len]
offset = end_offset - total_len
else:
data = video_bytes[offset:end_offset]
offset = end_offset
# 封装数据并发送到 NPU(硬件 VDEC 解码,不占用 CPU)
data_input = MxDataInput()
data_input.data = data
ret = stream_manager.SendData(stream_name, in_plugin_id, data_input)
if ret != 0:
error_count += 1
print(f"{prefix} SendData 失败, ret={ret} (累计失败: {error_count})")
if error_count >= 10:
print(f"{prefix} 连续失败次数过多,终止本路压测。")
break
continue
error_count = 0 # 发送成功则重置错误计数
# 从 NPU 拉取推理结果(若当前帧尚未解码完成,此处会短暂阻塞)
result = stream_manager.GetResult(stream_name, 0)
if result.errorCode == 0:
frame_count += 1
# 5. 统计并打印压测结果
actual_duration = time.time() - start_time
fps = frame_count / actual_duration if actual_duration > 0 else 0
print(
f"{prefix} 压测结束 | "
f"耗时: {actual_duration:.2f}s | "
f"获取结果次数: {frame_count} | "
f"处理频率: {fps:.2f} 次/秒"
)
# 6. 销毁资源
stream_manager.DestroyAllStreams()
def main():
parser = argparse.ArgumentParser(
description="NPU VDEC(硬件解码) + 推理并发性能压测工具",
formatter_class=argparse.RawTextHelpFormatter
)
parser.add_argument("--concurrency", type=int, default=10,
help="并发路数 (默认: 10)")
parser.add_argument("--video", type=str, required=True,
help="测试用的 H.264 裸流文件路径 (.h264)\n"
"抽取命令: ffmpeg -i input.mp4 -codec copy -bsf:v h264_mp4toannexb output.h264")
parser.add_argument("--model", type=str, required=True,
help="测试用的 OM 模型路径 (.om)")
parser.add_argument("--device", type=int, default=0,
help="NPU 设备 ID (默认: 0)")
parser.add_argument("--duration", type=int, default=60,
help="测试持续时间(秒) (默认: 60)")
args = parser.parse_args()
# ----------------------------------------------------------------
# 【关键修复】在主进程中提前解析绝对路径
# 使用 os.path.realpath() 而非 abspath(),可额外解析软链接
# 子进程直接使用此路径,避免 fork 后 CWD 不一致导致路径解析失败
# ----------------------------------------------------------------
abs_model_path = os.path.realpath(args.model)
abs_video_path = os.path.realpath(args.video)
# 启动前做完整的前置检查,避免所有子进程都失败才发现问题
print("=== 前置检查 ===")
has_error = False
if not os.path.isfile(abs_model_path):
print(f"[Error] 模型文件不存在: {abs_model_path}")
has_error = True
elif not os.access(abs_model_path, os.R_OK):
print(f"[Error] 模型文件无读取权限: {abs_model_path}")
has_error = True
else:
print(f"[OK] 模型路径: {abs_model_path}")
if not os.path.isfile(abs_video_path):
print(f"[Error] 视频文件不存在: {abs_video_path}")
has_error = True
elif not os.access(abs_video_path, os.R_OK):
print(f"[Error] 视频文件无读取权限: {abs_video_path}")
has_error = True
else:
video_size_mb = os.path.getsize(abs_video_path) / (1024 * 1024)
print(f"[OK] 视频路径: {abs_video_path} ({video_size_mb:.2f} MB)")
if has_error:
sys.exit(1)
print(f"\n=== NPU 压测启动 ===")
print(f"并发路数 : {args.concurrency}")
print(f"NPU 设备 : {args.device}")
print(f"持续时间 : {args.duration} 秒")
print("====================\n")
# 启动多进程并发压测
processes = []
for i in range(args.concurrency):
p = multiprocessing.Process(
target=run_single_stream,
args=(i, abs_video_path, abs_model_path, args.device, args.duration)
)
processes.append(p)
p.start()
# 主进程提示监控命令
print("=" * 55)
print("👀 压测正在后台进行,请打开另一个终端执行监控命令:")
print(f" watch -n 1 'npu-smi info -t usage -i {args.device} -c 0'")
print(" 重点观察:VDEC 利用率 / AI Core 利用率 / Memory Usage")
print("=" * 55 + "\n")
for p in processes:
p.join()
print("\n=== 压测全部完成 ===")
if __name__ == "__main__":
main()
Mind sdk版本5.0,CANN版本7.0。一直测试不通过,有大佬碰到过吗
MxpiVideoDecoder.cpp:114] The width of video is out of range [128,4096], actual width is 0. (Code = 1009, Message = "out of range") (Code = 1009, Message = "out of range")
import os import sys import time import multiprocessing import argparse import json try: # 导入华为 MindX SDK (mxVision) from StreamManagerApi import StreamManagerApi, MxDataInput except ImportError: print("[Error] 无法导入 StreamManagerApi 或 MxDataInput。") print("要在 Python 中测试 VDEC 硬件解码,必须正确安装华为的 MindX SDK (mxVision)。") print("请确认服务器的环境变量是否已加载 (source set_env.sh),或检查 libffi 动态库冲突。") sys.exit(1) def build_pipeline_config(device_id, model_path, stream_id): pipeline_name = f"vdec_infer_pipeline_{stream_id}" pipeline = { pipeline_name: { "stream_config": { "deviceId": str(device_id) }, "appsrc0": { "factory": "appsrc", "props": { "blocksize": "409600" }, "next": "h264parse0" }, "h264parse0": { "factory": "h264parse", "next": "mxpi_videodecoder0" }, # mxpi_videodecoder0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0113.html "mxpi_videodecoder0": { "factory": "mxpi_videodecoder", "props": {}, "next": "mxpi_imageresize0" }, # mxpi_imageresize0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0111.html "mxpi_imageresize0": { "factory": "mxpi_imageresize", "props": { "resizeType": "Resizer_Stretch", "resizeHeight": "640", "resizeWidth": "640" }, "next": "mxpi_tensorinfer0" }, # mxpi_tensorinfer0对应的文档地址https://www.hiascend.com/document/detail/zh/mindsdk/500/vision/visionug/mxvisionug_0122.html "mxpi_tensorinfer0": { "factory": "mxpi_tensorinfer", "props": { "dataSource": "mxpi_imageresize0", "modelPath": model_path }, "next": "appsink0" }, "appsink0": { "factory": "appsink" } } } return json.dumps(pipeline).encode('utf-8'), pipeline_name.encode('utf-8') def run_single_stream(stream_id, video_path, model_path, device_id, test_duration=60): """ 单路视频流的压测逻辑(运行在独立子进程中)。 参数: stream_id : 当前流的编号,用于日志区分和唯一 stream_name 生成 video_path : H.264 裸流文件的绝对路径(由主进程传入) model_path : OM 模型文件的绝对路径(由主进程传入) device_id : NPU 设备 ID test_duration : 压测持续时间(秒) """ prefix = f"[Stream {stream_id}]" # 1. 初始化 StreamManager stream_manager = StreamManagerApi() ret = stream_manager.InitManager() if ret != 0: print(f"{prefix} 初始化 StreamManager 失败, ret={ret}") return # 2. 构建并创建 Pipeline 流 pipeline_conf, stream_name = build_pipeline_config(device_id, model_path, stream_id) ret = stream_manager.CreateMultipleStreams(pipeline_conf) if ret != 0: print(f"{prefix} 创建 Pipeline 失败, ret={ret}") print(f"{prefix} 请确认模型路径是否正确: {model_path}") stream_manager.DestroyAllStreams() return in_plugin_id = 0 print(f"{prefix} Pipeline 创建成功,开始压测... (stream_name={stream_name.decode()})") # 3. 读取本地 H.264 裸流文件 # 注意:必须是用 ffmpeg 抽取的 .h264 裸流,不能直接使用 .mp4 容器格式 # 抽取命令示例:ffmpeg -i input.mp4 -codec copy -bsf:v h264_mp4toannexb output.h264 try: with open(video_path, "rb") as f: video_bytes = f.read() except Exception as e: print(f"{prefix} 视频文件读取失败: {e}") stream_manager.DestroyAllStreams() return if len(video_bytes) == 0: print(f"{prefix} 视频文件为空,退出。") stream_manager.DestroyAllStreams() return start_time = time.time() frame_count = 0 error_count = 0 chunk_size = 409600 # 必须与 pipeline 中 appsrc0 的 blocksize 匹配 (400KB) total_len = len(video_bytes) offset = 0 # 4. 循环向 NPU 灌入数据并拉取推理结果 while time.time() - start_time < test_duration: # 循环读取视频字节(到末尾后从头开始,模拟持续流) end_offset = offset + chunk_size if end_offset > total_len: data = video_bytes[offset:] + video_bytes[:end_offset - total_len] offset = end_offset - total_len else: data = video_bytes[offset:end_offset] offset = end_offset # 封装数据并发送到 NPU(硬件 VDEC 解码,不占用 CPU) data_input = MxDataInput() data_input.data = data ret = stream_manager.SendData(stream_name, in_plugin_id, data_input) if ret != 0: error_count += 1 print(f"{prefix} SendData 失败, ret={ret} (累计失败: {error_count})") if error_count >= 10: print(f"{prefix} 连续失败次数过多,终止本路压测。") break continue error_count = 0 # 发送成功则重置错误计数 # 从 NPU 拉取推理结果(若当前帧尚未解码完成,此处会短暂阻塞) result = stream_manager.GetResult(stream_name, 0) if result.errorCode == 0: frame_count += 1 # 5. 统计并打印压测结果 actual_duration = time.time() - start_time fps = frame_count / actual_duration if actual_duration > 0 else 0 print( f"{prefix} 压测结束 | " f"耗时: {actual_duration:.2f}s | " f"获取结果次数: {frame_count} | " f"处理频率: {fps:.2f} 次/秒" ) # 6. 销毁资源 stream_manager.DestroyAllStreams() def main(): parser = argparse.ArgumentParser( description="NPU VDEC(硬件解码) + 推理并发性能压测工具", formatter_class=argparse.RawTextHelpFormatter ) parser.add_argument("--concurrency", type=int, default=10, help="并发路数 (默认: 10)") parser.add_argument("--video", type=str, required=True, help="测试用的 H.264 裸流文件路径 (.h264)\n" "抽取命令: ffmpeg -i input.mp4 -codec copy -bsf:v h264_mp4toannexb output.h264") parser.add_argument("--model", type=str, required=True, help="测试用的 OM 模型路径 (.om)") parser.add_argument("--device", type=int, default=0, help="NPU 设备 ID (默认: 0)") parser.add_argument("--duration", type=int, default=60, help="测试持续时间(秒) (默认: 60)") args = parser.parse_args() # ---------------------------------------------------------------- # 【关键修复】在主进程中提前解析绝对路径 # 使用 os.path.realpath() 而非 abspath(),可额外解析软链接 # 子进程直接使用此路径,避免 fork 后 CWD 不一致导致路径解析失败 # ---------------------------------------------------------------- abs_model_path = os.path.realpath(args.model) abs_video_path = os.path.realpath(args.video) # 启动前做完整的前置检查,避免所有子进程都失败才发现问题 print("=== 前置检查 ===") has_error = False if not os.path.isfile(abs_model_path): print(f"[Error] 模型文件不存在: {abs_model_path}") has_error = True elif not os.access(abs_model_path, os.R_OK): print(f"[Error] 模型文件无读取权限: {abs_model_path}") has_error = True else: print(f"[OK] 模型路径: {abs_model_path}") if not os.path.isfile(abs_video_path): print(f"[Error] 视频文件不存在: {abs_video_path}") has_error = True elif not os.access(abs_video_path, os.R_OK): print(f"[Error] 视频文件无读取权限: {abs_video_path}") has_error = True else: video_size_mb = os.path.getsize(abs_video_path) / (1024 * 1024) print(f"[OK] 视频路径: {abs_video_path} ({video_size_mb:.2f} MB)") if has_error: sys.exit(1) print(f"\n=== NPU 压测启动 ===") print(f"并发路数 : {args.concurrency}") print(f"NPU 设备 : {args.device}") print(f"持续时间 : {args.duration} 秒") print("====================\n") # 启动多进程并发压测 processes = [] for i in range(args.concurrency): p = multiprocessing.Process( target=run_single_stream, args=(i, abs_video_path, abs_model_path, args.device, args.duration) ) processes.append(p) p.start() # 主进程提示监控命令 print("=" * 55) print("👀 压测正在后台进行,请打开另一个终端执行监控命令:") print(f" watch -n 1 'npu-smi info -t usage -i {args.device} -c 0'") print(" 重点观察:VDEC 利用率 / AI Core 利用率 / Memory Usage") print("=" * 55 + "\n") for p in processes: p.join() print("\n=== 压测全部完成 ===") if __name__ == "__main__": main()