【Scala PyTorch深度学习】PyTorch On Scala 系列课程 第十二章 25 :PyTorch算子模型优化【AI Infra 3.0】[PyTorch Scala 硕士研一课程]

PyTorch Scala 高校计算机硕士研一课程
通过外部库优化算子
尽管 PyTorch 通过其 ATen 后端提供了大量优化的操作库,但有时模型或数据处理管道中特定、自定义计算步骤会出现性能瓶颈。这些瓶颈可能源于复杂的逐元素操作、未能有效映射到标准 PyTorch 函数的算法,或者需要对 GPU 执行进行精细控制。当 PyTorch 分析器识别出此类算子是性能限制因素时,使用专门用于加速数值计算的外部库会是一种有效的优化方法。
本节研究如何将 CuPy 和 Numba 等库整合到您的 PyTorch 工作流程中,以加速这些重要的计算算子,补充本章讨论的更广泛的部署优化技术。
使用 CuPy 加速 GPU 计算
CuPy 是一个开源库,它提供了一个与 NumPy 兼容的多维数组接口,并使用 NVIDIA CUDA 进行加速。如果您的瓶颈涉及 GPU 上复杂的数组操作,而这些操作可能更自然地用 NumPy 风格的索引和操作来表示,或者如果您需要在不承担构建 C++ 扩展(第 6 章讨论)的全部开销的情况下编写自定义 CUDA 算子,CuPy 是一个有力的选择。
整合 CuPy 与 PyTorch
其主要思路是将张量数据从 PyTorch 传输到 CuPy,使用 CuPy 的函数或自定义算子执行加速计算,然后将结果传回 PyTorch。PyTorch 和 CuPy 的现代版本支持 DLPack 标准,这允许在同一设备上的库之间进行零拷贝数据共享,从而显著减少开销。
- PyTorch 张量到 CuPy 数组: 您可以使用
cupy.asarray()将 PyTorch GPU 张量转换为 CuPy 数组。如果支持 DLPack 并且张量位于同一 GPU 设备上,此操作通常可以避免数据拷贝。 - CuPy 计算: 使用 CuPy 丰富的函数集执行计算,这些函数模仿 NumPy 的 API 但在 GPU 上执行。您还可以使用 CuPy 的
cupy.RawKernel定义并启动自定义 CUDA 算子。 - CuPy 数组到 PyTorch 张量: 使用
torch.as_tensor()将生成的 CuPy 数组转换回 PyTorch 张量。同样,如果数组位于 PyTorch 识别的 CUDA 设备上,DLPack 有助于高效、可能零拷贝的传输。
示例:使用 CuPy 进行自定义逐元素操作
设想一个在纯 Python 或标准 PyTorch 操作中表现缓慢的自定义激活函数:
import torch
import cupy
import math
// 使用 CuPy 的逐元素核函数特性定义自定义操作
// 示例:如果 x < threshold 则 y = log(1 + exp(x)) 否则 y = x
val custom_softplus_kernel = cupy.ElementwiseKernel(
'T x, float64 threshold', # 输入参数
'T y', # 输出参数
'''
if (x < threshold) {
y = log(1.0 + exp(x));
} else {
y = x;
}
''',
'custom_softplus' # 核函数名称
)
// GPU 上的 PyTorch 示例张量
val pytorch_tensor_gpu = torch.randn(1000, 1000, device='cuda')
// 1. 将 PyTorch 张量转换为 CuPy 数组(可能通过 DLPack 实现零拷贝)
val cupy_array = cupy.asarray(pytorch_tensor_gpu)
// 2. 应用自定义 CuPy 核函数
val threshold_value = 10.0
val result_cupy_array = custom_softplus_kernel(cupy_array, threshold_value)
// 3. 将结果转换回 PyTorch 张量(可能通过 DLPack 实现零拷贝)
val result_pytorch_tensor = torch.as_tensor(result_cupy_array, device='cuda')
// 如果需要进行计时或后续 CPU 操作,请确保同步
// torch.cuda.synchronize()
println(s"输入张量设备: ${pytorch_tensor_gpu.device}")
println(s"结果张量设备: ${result_pytorch_tensor.device}")
println(s"结果张量形状: ${result_pytorch_tensor.shape}")
何时使用 CuPy:
- 您的瓶颈涉及复杂数组操作,这些操作易于用 NumPy 语法表达但需要 GPU 加速。
- 您需要编写中等复杂度的自定义 CUDA 算子,而无需设置完整的 C++/CUDA 扩展构建系统。
- 核函数的计算成本足够高,可以抵消 PyTorch-CuPy 数据接口带来的任何潜在开销。
使用 Numba 进行即时编译
Numba 是另一个功能强大的库,它使用 LLVM 编译器基础设施在运行时将 Python 函数转换为优化的机器代码。它既可以针对 CPU,也可以针对 NVIDIA GPU(通过 numba.cuda 子模块)。与提供 CUDA 加速的 NumPy 替代方案的 CuPy 不同,Numba 侧重于加速您已有的 Python 代码,通常只需进行最少的修改(例如添加装饰器)。
将 Numba 与 PyTorch 数据配合使用
Numba 不直接操作 PyTorch 张量。您通常需要:
- 访问底层数据,通常通过将 PyTorch 张量转换为 NumPy 数组(对于 CPU 操作使用
tensor.cpu().numpy(),如果 Numba 针对 CUDA 则可能使用 DLPack/CuPy 作为中间层来访问 GPU 数据)。 - 对此数据应用 Numba 装饰的函数(
@numba.jit、@numba.vectorize或@numba.cuda.jit)。 - 如有必要,将结果转换回 PyTorch 张量。
示例:使用 Numba JIT 进行 CPU 密集型计算
假设您在 CPU 上有一个复杂的后处理步骤,其中涉及在纯 Python 中运行缓慢的循环。
import torch
import numpy as np
import numba
// 定义一个可能在 NumPy 数组上运行缓慢的 Python 函数
@numba.jit(nopython=True) // 使用 nopython=True 以获得最佳性能
def complex_cpu_calculation(data_array, scale_factor):
rows, cols = data_array.shape
result = np.empty_like(data_array)
for i in range(rows):
for j in range(cols):
val = data_array[i, j]
# 复杂计算示例
processed_val = (np.sin(val) * scale_factor + np.cos(val / scale_factor))**2
result[i, j] = processed_val
return result
// CPU 上的 PyTorch 示例张量
val pytorch_tensor_cpu = torch.randn(500, 500, device='cpu')
// 1. 转换为 NumPy 数组(CPU 张量零拷贝)
val numpy_array = pytorch_tensor_cpu.numpy()
// 2. 应用 Numba 加速函数
val scale = 2.5
val result_numpy_array = complex_cpu_calculation(numpy_array, scale)
// 3. 转换回 PyTorch 张量(CPU 张量零拷贝)
val result_pytorch_tensor = torch.from_numpy(result_numpy_array)
println(s"输入张量设备: ${pytorch_tensor_cpu.device}")
println(s"结果张量设备: ${result_pytorch_tensor.device}")
println(s"结果张量形状: ${result_pytorch_tensor.shape}")
将 Numba 用于 CUDA 算子
Numba 还允许使用 @numba.cuda.jit 直接以 Python 语法编写 CUDA 算子。对于不太复杂的 GPU 任务,这可能比 CuPy 的 RawKernel 或完整的 C++ 扩展更简单。
import torch
import numpy as np
import numba
import numba.cuda
import math
@cuda.jit
def gpu_kernel(x, out):
idx = cuda.grid(1) // 获取全局线程索引
if idx < x.shape[0]:
# 逐元素 GPU 操作示例
out[idx] = math.exp(math.sin(x[idx]) * 2.0)
// GPU 上的 PyTorch 示例张量
val pytorch_tensor_gpu = torch.randn(2**16, device='cuda')
// Numba CUDA 要求支持 CUDA 数组接口的类数组对象
// 最简单的方法通常是通过 NumPy/CuPy 中间件,或者如果兼容则直接访问
// 注意:直接使用 pytorch_tensor_gpu.__cuda_array_interface__ 可能有效
// 但明确使用 CuPy 通常能使 GPU 到 Numba 的交互更清晰。
// 使用 CuPy 作为中间件(推荐,以提高清晰度)
import cupy
val cupy_array_in = cupy.asarray(pytorch_tensor_gpu)
val cupy_array_out = cupy.empty_like(cupy_array_in)
// 配置线程/块维度
val threads_per_block = 128
val blocks_per_grid = (cupy_array_in.size + (threads_per_block - 1)) // threads_per_block
// 启动 Numba CUDA 核函数
gpu_kernel[blocks_per_grid, threads_per_block](cupy_array_in, cupy_array_out)
// 将结果转换回 PyTorch 张量
val result_pytorch_tensor = torch.as_tensor(cupy_array_out, device='cuda')
println(s"输入张量设备: ${pytorch_tensor_gpu.device}")
println(s"结果张量设备: ${result_pytorch_tensor.device}")
println(s"结果张量形状: ${result_pytorch_tensor.shape}")
何时使用 Numba:
- 您的瓶颈在于纯 Python 代码(循环、复杂逻辑),操作 NumPy 兼容数据,无论是在 CPU 还是 GPU 上。
- 您更喜欢主要通过 Python 语法并使用装饰器来编写优化代码。
@numba.jit(nopython=True)模式适用于显著的 CPU 加速。- 您需要编写相对简单的自定义 CUDA 算子,而无需外部编译步骤(
@numba.cuda.jit)。
权衡与考量
整合 CuPy 或 Numba 等外部库能带来潜在的性能提升,但也引入了需要考虑的因素:
- 数据传输开销: 在 PyTorch 和这些库之间移动数据(即使有 DLPack 等零拷贝机制)也会有一定开销。确保优化后的算子内部的计算量足够大,以证明此成本是合理的。小型操作的频繁往返可能导致性能下降。
- 依赖性: 为您的项目添加 CuPy 或 Numba 作为依赖项,可能使部署环境变得复杂。
- 复杂性: 引入了另一层抽象,并需要理解所选库的细节(CuPy API、Numba 编译模式、用于 GPU 的 CUDA 知识)。
- 调试: 调试跨多个库(PyTorch、CuPy/Numba,可能还有 CUDA)的代码更具挑战性。
使用外部库优化特定算子是一种有针对性的方法。在分析已识别出标准 PyTorch 操作或其他优化技术(如 TorchScript 或量化)无法充分解决的明确的、计算密集型瓶颈之后,这种方法最有效。通过审慎地整合 CuPy 和 Numba 等工具,您可以显著加速那些重要部分,从而有助于模型部署得更快、效率更高。
模型导出为 ONNX 格式
收藏
虽然 TorchScript 在 PyTorch 生态系统中提供了一种序列化 PyTorch 模型的方法,但为了实现更广泛的互操作性,通常需要一种标准化格式。开放神经网络交换 (ONNX) 格式满足了这一要求,它定义了一个开放标准来表示机器学习模型。将您的 PyTorch 模型导出为 ONNX 可以使它们在多种平台和推理引擎上运行,例如 ONNX Runtime、TensorRT、OpenVINO 以及各种移动/边缘设备,并且通常可以从这些运行时提供的硬件特定优化中获益。将 PyTorch 模型转换为 ONNX 格式的过程进行了详细说明。
ONNX 在部署中的作用
ONNX 充当中间表示形式。您可以使用 PyTorch 灵活的环境训练模型,然后将训练好的模型图及其学习到的参数导出到 .onnx 文件。此文件随后可以由任何 ONNX 兼容的运行时加载和执行。这种解耦显著简化了部署过程,因为您不需要在每个目标部署系统上都安装 PyTorch。它还通过允许专用运行时应用图优化并更有效地使用加速器,从而提升了性能,这可能比通用框架更有效。
推理环境PyTorch 训练torch.onnx.export()Model.onnxONNX Runtime加载并运行NVIDIA TensorRT加载并运行Intel OpenVINO加载并运行Apple Core ML加载并运行其他运行时…加载并运行
工作流程示意 PyTorch 模型导出到 ONNX 以及随后在各种推理运行时中的部署。
使用 torch.onnx.export 导出模型
PyTorch 在 torch.onnx 模块中提供了 torch.onnx.export() 函数作为此转换的主要工具。此函数的核心功能通常是利用追踪来记录当样本输入通过模型时执行的操作,并将这些操作转换为其 ONNX 等效项。
该函数签名有几个重要参数:
torch.onnx.export(
model, // 要导出的模型 (torch.nn.Module)
args, // 用于追踪的模型输入元组
f, // 输出路径(字符串)或类文件对象
export_params=true, // 在文件中存储训练过的参数
opset_version=None, // ONNX 算子集版本
do_constant_folding=true, // 执行常量折叠优化
input_names=None, // ONNX 图中输入节点名称列表
output_names=None, // ONNX 图中输出节点名称列表
dynamic_axes=None // 指定动态维度的字典
// ... 其他参数
)
参数
model:您的torch.nn.Module实例。如果模型在训练和推理之间的行为(例如 dropout 或批归一化)有所不同,请确保其处于评估模式 (model.eval())。args:一个元组,包含具有正确数据类型和形状的示例输入,这些是您的模型forward方法所期望的。此输入用于追踪执行路径。重要地,args中的形状定义了导出的 ONNX 图中的输入形状,除非使用了dynamic_axes。f:.onnx模型将要保存的文件路径。export_params:如果为True(默认),模型的训练权重将直接嵌入到 ONNX 文件中,使其成为自包含文件。opset_version:指定要使用的 ONNX 算子集版本。不同的版本支持不同的算子集和功能。选择正确的 opset 对于与目标推理运行时的兼容性很重要。请查阅目标运行时的文档以获取支持的 opset。常见选择介于 11 到 17 之间,但新版本会定期发布。input_names/output_names:可选的字符串列表,为 ONNX 图中的输入和输出节点提供有意义的名称。这提高了可读性,并使得在运行时使用 ONNX 模型时更容易提供数据和获取结果。dynamic_axes:这是处理可变输入/输出形状的一个非常重要的参数。
处理动态形状
追踪本质上会捕获所提供 args 的特定形状。如果您的模型需要处理不同维度的输入(例如,NLP 模型中的可变批大小或序列长度),您必须使用 dynamic_axes 参数来指定这一点。
dynamic_axes 是一个字典,其中键是前面定义的 input_names 或 output_names,值是另一个字典,将轴索引映射到描述性名称。例如,要指定名为 ‘input_ids’ 的输入的批大小(轴 0)和序列长度(轴 1)可以变化,您可以使用:
dynamic_axes = {
'input_ids': {0: 'batch_size', 1: 'sequence_length'}, # 输入轴定义
'output_logits': {0: 'batch_size'} # 输出轴定义
}
这会告诉导出器不要将这些维度硬编码到图中,从而允许 ONNX 运行时处理沿这些指定轴具有不同大小的输入并生成输出。
导出实践示例
我们来导出一个简单的卷积模型。
import torch
import torch.nn as nn
import torch.onnx
// 定义一个简单的 CNN 模型
class SimpleCNN extends nn.Module:
def __init__(self):
super().__init__()
val conv1 = nn.Conv2d(3, 16, kernel_size=3, stride=1, padding=1)
val relu = nn.ReLU()
val pool = nn.MaxPool2d(kernel_size=2, stride=2)
val fc1 = nn.Linear(16 * 16 * 16, 10) // 假设输入图像为 32x32
def forward(x: torch.Tensor):
x = pool(relu(conv1(x)))
x = torch.flatten(x, 1) // 展平除批次维度外的所有维度
x = fc1(x)
return x
// 实例化模型并设置为评估模式
val model = SimpleCNN()
model.eval()
// 创建与预期维度匹配的虚拟输入(批大小、通道、高度、宽度)
// 注意:这里批大小设置为 1,但我们会使其变为动态
val dummy_input = torch.randn(1, 3, 32, 32, requires_grad=False)
// 定义输入和输出名称
val input_names = List("input_image")
val output_names = List("output_logits")
// 定义动态轴(使批大小动态化)
val dynamic_axes_config = Map(
"input_image" -> Map(0 -> "batch_size"), // 输入的可变批大小
"output_logits" -> Map(0 -> "batch_size") // 输出的可变批大小
)
// 指定输出文件路径
val onnx_model_path = "simple_cnn.onnx"
// 导出模型
torch.onnx.export(
model,
dummy_input,
onnx_model_path,
export_params=true,
opset_version=12, // 选择合适的 opset 版本
do_constant_folding=true,
input_names=input_names,
output_names=output_names,
dynamic_axes=dynamic_axes_config
)
println(s"模型已导出到 $onnx_model_path")
常见的导出挑战
虽然 torch.onnx.export 对许多模型都适用,但您可能会遇到一些问题:
- 不支持的 PyTorch 算子: 并非每个 PyTorch 函数或模块在目标 ONNX
opset_version中都有直接的等效项。如果追踪器遇到不支持的操作,导出将会失败。解决方案包括:- 重写 PyTorch 代码以使用 ONNX 兼容的操作。
- 选择一个可能支持该算子的不同(通常是更新的)
opset_version。 - 实现自定义 ONNX 算子(一个涉及 C++ 的高级主题)。
- 如果操作仅在训练期间发生,请确保模型处于
eval()模式。
- 动态控制流: 追踪难以处理依赖于数据的控制流(例如,条件或迭代次数取决于张量值的
if语句或循环)。虽然torch.jit.script有时可以捕获此类逻辑,但将脚本化模型导出到 ONNX 也可能具有挑战性。通常需要简化控制流或使其与数据无关。 - Opset 兼容性: 导出的 ONNX 模型必须使用目标推理引擎(例如 ONNX Runtime)支持的 opset 版本。请务必查看运行时的文档以了解兼容的 opset。
验证导出模型
导出后,验证 ONNX 模型的正确性很重要。一种常见的方法是使用 onnxruntime 库:
import onnxruntime as ort
import numpy as np
// 加载 ONNX 模型
val ort_session = ort.InferenceSession(onnx_model_path)
// 准备输入数据(需要是 NumPy 数组)
// 创建一个不同批大小的输入以测试动态轴
val test_input_np = np.random.randn(4, 3, 32, 32).astype(np.float32) // 批大小 = 4
// 运行推理
val ort_inputs = Map(ort_session.get_inputs().head.name -> test_input_np)
val ort_outputs = ort_session.run(None, ort_inputs)
val onnx_result = ort_outputs.head
// 与 PyTorch 输出进行比较(可选,但建议)
// 如果 ONNX Runtime 使用 CPU,请确保模型在 CPU 上以便直接比较
model.cpu()
val dummy_input = torch.from_numpy(test_input_np)
with torch.no_grad():
val pytorch_result = model(dummy_input).numpy()
// 检查输出是否接近(允许存在潜在的微小数值差异)
if np.allclose(pytorch_result, onnx_result, rtol=1e-03, atol=1e-05):
println("验证成功:ONNX Runtime 输出与 PyTorch 输出一致。")
else:
println("验证失败:输出不一致。")
// 可能需要进一步调试
此验证步骤有助于确保转换过程没有引入错误,并且模型在目标运行时环境中表现符合预期,至少在数值上是这样。
导出到 ONNX 是一种有用的技术,可以使您的先进 PyTorch 模型可移植,并为在各种硬件和软件平台上的高效部署做好准备。掌握这一过程,包括处理动态形状和解决常见问题,是使您的深度学习应用程序投入生产的重要一步。
使用 TorchServe 提供模型服务
收藏
将机器学习模型部署到最终用户或下游应用程序,需要通过服务基础设施使其可用。这通常发生在模型使用 TorchScript、量化或剪枝等技术优化推理性能之后。手动构建和管理此类服务基础设施可能很复杂,涉及 API 开发、请求处理、扩缩容和监控。TorchServe 是一个专门为简化 PyTorch 模型这一过程而开发的工具。它提供了一种标准化的方式,用于在生产环境中打包、部署、管理和提供训练好的模型。
TorchServe 充当了桥梁,连接优化后的 PyTorch 模型与需要使用其预测结果的应用程序。它处理模型服务的操作方面,使您能够专注于模型开发和集成。
理解 TorchServe 架构
TorchServe 在设计时充分考虑了灵活性和性能。其主要组件共同协作以提供一个方案:
- 模型归档器 (
torch-model-archiver): 这个命令行工具是为 TorchServe 准备模型的第一步。它将所有必需的工件打包成一个独立的归档文件,扩展名为.mar。这些工件通常包括:- 序列化模型文件(例如,TorchScript
.pt文件或标准的state_dict)。 - 定义 处理器 逻辑(预处理、推理、后处理)的 Python 脚本。
- 可选的辅助文件(例如,词汇文件、配置 JSON、标签映射文件)。
- 清单文件(自动生成),描述模型、版本、处理器等。
- 序列化模型文件(例如,TorchScript
- TorchServe 运行时: 这是核心服务器进程。它侦听预定义网络端口上的传入请求。它管理已部署模型的生命周期,包括将它们加载到内存中、按模型扩展工作进程数量,以及将推理请求路由到适当的工作进程。
- 处理器: 处理器是 Python 脚本或类,它们定义 TorchServe 如何与您的特定模型交互。它们封装了以下逻辑:
initialize(context): 模型加载时调用一次。用于将模型加载到内存和进行一次性设置。preprocess(data): 将传入请求数据转换为模型forward方法期望的格式(例如,解码图像、文本分词)。inference(model_input): 通过调用模型的预测函数来执行实际推理。postprocess(inference_output): 将模型的原始输出转换为用户友好的格式(例如,将类别索引映射到标签,格式化 JSON)。 TorchServe 提供了用于常见任务(图像分类、对象检测、文本分类)的多个内置处理器,但您可以轻松地为特殊的模型输入/输出或复杂的工作流创建自定义处理器。
- API: TorchServe 公开了两个主要 REST API:
- 推理 API (默认端口: 8080): 用于向已加载的模型发送推理请求并接收预测结果。
- 管理 API (默认端口: 8081): 用于管理 TorchServe 提供的模型。操作包括注册新模型、注销现有模型、设置模型的默认版本以及按模型扩展工作进程数量。
- 指标 API (默认端口: 8082): 以 Prometheus 兼容的格式公开有关服务器和模型的运行指标(例如,请求延迟、错误率、CPU/内存使用情况)。
gRPC 支持也可用,用于低延迟通信,在微服务架构中尤其有用。
The TorchServe 工作流程
使用 TorchServe 部署模型通常遵循以下步骤:
- 准备模型工件: 确保您的训练模型已保存(例如,使用
torch.jit.save()保存 TorchScript 模型或使用torch.save()保存state_dict),并收集所有必需的辅助文件。 - 编写处理器(如果需要): 如果内置处理器不适合您的模型,请实现一个自定义处理器脚本(
.py文件)。 - 归档模型: 使用
torch-model-archiver将模型、处理器和其他文件打包成一个.mar归档。 - 启动 TorchServe: 启动 TorchServe 运行时,将其指向一个目录(“模型存储”),
.mar文件将位于或注册到该目录。 - 注册模型: 使用管理 API 告知 TorchServe 从其
.mar文件加载您的模型并准备好服务。您可以指定初始工作进程数量。 - 发送推理请求: 使用推理 API 将数据发送到注册模型的端点并接收预测结果。
- 管理与监控: 使用管理 API 扩缩工作进程或更新模型,并使用指标 API 监控性能。
准备服务运行时客户端应用程序管理 API(端口 8081)注册/扩缩/注销推理 API(端口 8080)推理请求TorchServe(前端 / 运行时)处理器(初始化, 预处理,推理, 后处理)路由请求后处理结果PyTorch 模型(.pt / state_dict)预处理数据torch-model-archiver原始输出模型归档(.mar 文件)创建模型存储(目录)放入 / 通过 API 注册从…加载预测结果指标 API(端口 8082)监控系统(例如,Prometheus)抓取
TorchServe 部署工作流程的高级概述,描绘了准备步骤和运行时请求处理。
创建模型归档
torch-model-archiver 工具对于模型打包非常核心。以下是典型的命令结构:
torch-model-archiver \
--model-name my_transformer_model \
--version 1.0 \
--serialized-file traced_transformer.pt \
--handler transformer_handler.py \
--extra-files "vocab.txt,config.json" \
--export-path /path/to/model-store \
--force
让我们分析一下这些参数:
--model-name: 在 API 调用中使用的模型逻辑名称(例如,my_transformer_model)。--version: 模型的版本字符串(例如,1.0)。TorchServe 可以管理同一模型的多个版本。--serialized-file: 已保存模型文件的路径(例如,torch.jit.save的输出)。如果使用state_dict,您还需要--model-file指向您的模型定义 Python 文件。--handler: 您的处理器脚本(.py)的路径。这可以是内置处理器之一(例如,image_classifier,text_classifier),也可以是您的自定义脚本。--extra-files: 模型或处理器所需的额外文件的逗号分隔列表(例如,分词器配置、词汇文件、标签映射)。这些文件将通过上下文对象在处理器的initialize方法中访问。--export-path: 生成的.mar文件将保存到的目录。这通常设置为 TorchServe 模型存储目录。--force: 如果.mar文件已存在,则覆盖它。
执行此命令将创建 /path/to/model-store/my_transformer_model.mar。
自定义处理器
虽然 TorchServe 的内置处理器涵盖了许多常见用例,但您通常会需要自定义逻辑。自定义处理器是一个 Python 脚本,其中包含一个类(通常继承自 BaseHandler),它实现了以下部分或全部方法:
// custom_handler.py
import torch
import json
from ts.torch_handler.base_handler import BaseHandler
import logging
import os
val logger = logging.getLogger(__name__)
class MyCustomHandler extends BaseHandler:
"""
自定义处理器,用于处理特定输入/输出格式。
"""
def __init__(self):
super().__init__()
val initialized = False
val model = None
val mapping = None
def initialize(self, context):
"""
加载模型和额外文件。模型加载时调用一次。
"""
val manifest = context.manifest
val properties = context.system_properties
val model_dir = properties.get("model_dir") // 包含解压后的 MAR 内容的目录
// 确定设备
val device = torch.device("cuda:" + str(properties.get("gpu_id")) if torch.cuda.is_available() and properties.get("gpu_id") is not None else "cpu")
logger.info(f"处理器在设备上初始化: {device}")
// 加载模型(假设为 TorchScript 的示例)
val serialized_file = manifest['model']['serializedFile']
val model_pt_path = os.path.join(model_dir, serialized_file)
if not os.path.isfile(model_pt_path):
raise RuntimeError("缺少 model.pt 文件")
val model = torch.jit.load(model_pt_path, map_location=device)
model.eval()
logger.info(f"模型 {manifest['model']['modelName']} 加载成功。")
// 加载额外文件(例如,标签映射)
val mapping_file_path = os.path.join(model_dir, "index_to_name.json") // 假设通过 --extra-files 传入
if os.path.isfile(mapping_file_path):
with open(mapping_file_path) as f:
self.mapping = json.load(f)
logger.info("标签映射加载成功。")
else:
logger.warning("映射文件未找到。")
val initialized = true
def preprocess(self, data):
"""
将原始输入数据转换为模型输入张量。
'data' 是一个字典列表,每个字典包含原始请求数据。
"""
// 示例:假设输入是 JSON,如 {'text': 'some input string'}
// 这在很大程度上取决于您的具体应用程序
val processed_inputs = []
for row in data:
val request_body = row.get("data") or row.get("body") // 处理不同的输入源
if isinstance(request_body, (bytes, bytearray)):
request_body = request_body.decode('utf-8')
// 在此处添加您的特定预处理逻辑(分词、张量创建等)
// 为简单起见,我们假设 request_body 是所需的直接输入
// 实际上,您会进行文本分词、图像解码/大小调整等操作。
// 这部分必须返回 self.model() 期望格式的数据。
logger.info(f"收到输入: {request_body}")
// 虚拟预处理:直接传递(替换为真实逻辑)
processed_inputs.append(request_body)
// 示例:如果模型需要,将处理后的输入转换为批处理张量
// input_tensor = self.tokenizer(processed_inputs, return_tensors="pt", padding=True, truncation=True).to(self.device)
// return input_tensor
return processed_inputs // 返回列表作为虚拟示例
def inference(self, model_input):
"""
使用模型运行推理。
'model_input' 是 preprocess() 的输出。
"""
// 示例:运行已加载的 TorchScript 模型
// 输出取决于您的模型结构
// 确保 model_input 在正确的设备上
// output = self.model(model_input.to(self.device))
// 虚拟推理:只回显输入(替换为真实模型调用)
logger.info(f"正在对: {model_input} 运行虚拟推理")
with torch.no_grad(): // 对推理来说必不可少
// 将此替换为:output = self.model(model_input)
output = [f"Processed: {item}" for item in model_input]
return output
def postprocess(self, inference_output):
"""
将模型输出转换为用户友好的格式。
'inference_output' 是 inference() 的输出。
返回预测结果列表,每个输入请求一个。
"""
// 示例:将原始模型输出(例如,logits)转换为标签/分数
// predictions = torch.softmax(inference_output, dim=1).argmax(dim=1).cpu().tolist()
// result = [self.mapping[str(pred)] if self.mapping else str(pred) for pred in predictions]
// 虚拟后处理:只返回推理输出
logger.info(f"后处理: {inference_output}")
return inference_output // 应该是一个列表
// 注意:处理器文件名必须与 --handler 参数值匹配(不带 .py)
// 如果处理器类名与首字母大写的文件名不同,
// 则在 MANIFEST.json 中或通过归档器指定。
此示例显示了基本结构。您将用实际的预处理(例如,图像转换、文本分词)、模型调用和后处理(例如,将输出索引映射到类别名称、格式化 JSON 响应)来替换虚拟逻辑。context 对象提供对系统属性(如 GPU 可用性)、模型目录和清单详情的访问。
运行 TorchServe 和管理模型
一旦您的 .mar 文件在模型存储中,您就可以启动 TorchServe:
# 如果模型存储目录不存在,则创建它
mkdir /path/to/model-store
# 启动 TorchServe,指向模型存储
torchserve --start \
--model-store /path/to/model-store \
--models my_model=/path/to/model-store/my_transformer_model.mar \
--ts-config /path/to/config.properties
--start: 在后台启动 TorchServe 服务器。使用torchserve --stop停止它。--model-store: 指定包含.mar文件的目录。--models:(可选)在启动时预加载并注册特定模型。格式为model_name=model_archive.mar,如果文件在模型存储中,也可以是简单的model_name.mar。您可以指定多个模型。--ts-config:(可选)用于自定义端口、JVM 参数、日志等的配置文件(.properties)路径。
启动后,您将与 API 交互,通常使用 curl 或 Python 中的 requests 等客户端库。
注册模型(管理 API):
curl -X POST "http://localhost:8081/models?url=my_transformer_model.mar&model_name=transformer&initial_workers=1&synchronous=true"
此命令告知 TorchServe 从其 .mar 文件加载 my_transformer_model.mar(假设它在模型存储中),将其注册为逻辑名称 transformer,为其启动 1 个工作进程,并等待注册完成。
发送推理请求(推理 API):
具体格式取决于您的处理器的 preprocess 方法。如果它期望原始图像字节:
curl -X POST http://localhost:8080/predictions/transformer -T image.jpg
如果它期望 JSON:
curl -X POST http://localhost:8080/predictions/transformer -H "Content-Type: application/json" -d '{"text": "This is an example sentence."}'
响应将包含来自您的处理器的 postprocess 方法的输出。
扩缩工作进程(管理 API):
如果您需要特定模型的更高吞吐量,您可以增加工作进程的数量:
curl -X PUT "http://localhost:8081/models/transformer?min_worker=4"
这会将 transformer 模型的工作进程数量扩缩到 4(TorchServe 处理它们之间的负载均衡)。
检查状态和指标:
- 列出已注册的模型:
curl http://localhost:8081/models - 描述特定模型:
curl http://localhost:8081/models/transformer - 获取指标:
curl http://localhost:8082/metrics
性能与可扩展性
TorchServe 专为生产负载而设计。有助于提高性能的特点包括:
- 工作进程: 在独立的工作进程中运行推理可以实现并行请求处理和隔离。
- 批处理: 处理器可以在
preprocess和inference中实现批处理逻辑,以同时处理多个请求,从而提高 GPU 利用率。TorchServe 还具有可通过管理 API 或配置文件配置的内置动态批处理能力。 - 异步后端: TorchServe 使用异步后端(基于 Netty)来高效处理并发连接。
- 指标端点: 提供详细指标,用于使用 Prometheus 和 Grafana 等工具监控性能并发现瓶颈。
- 集成: 可以轻松容器化 (Docker) 并使用 Kubernetes 等编排工具部署,以实现自动扩缩容和高可用性。
通过使用 TorchServe,您能显著减少部署和管理 PyTorch 模型所需的工程工作量,且确保可靠高效。它提供了一个标准化、功能丰富的平台,与 PyTorch 生态系统和常见的 MLOps 实践良好集成,使其成为将您的先进模型投入生产的重要工具。
实践:模型性能分析与量化
提供一个动手练习,将性能分析和量化方法应用于标准 PyTorch 模型。目标是使用性能分析工具确定性能特征。随后,通过训练后静态量化(PTQ)减小模型大小并可能加快其推理速度。此练习反映了准备模型部署时的常见流程。
我们假定您已安装 torchvision 并拥有可用的 PyTorch 环境。
环境和模型的准备
首先,我们需要导入所需的库并载入一个预训练模型。我们将使用 torchvision 中的 ResNet18 作为示例模型。它足够复杂,可显示有意义的结果,但又足够小,可在此练习中快速运行。我们还需要一些模拟输入数据。
import torch
import torchvision.models as models
import torch.quantization
import torch.profiler
import copy
import time
import os
import numpy as np
// 检查 CUDA 是否可用,如果不可用则回退到 CPU
val device = torch.device("cuda" if torch.cuda.is_available() else "cpu")
println(s"Using device: $device")
// 载入预训练的 ResNet18 模型
val original_model = models.resnet18(pretrained=True)
original_model.eval() // 将模型设置为评估模式
val model_fp32 = copy.deepcopy(original_model).to(device)
// 创建与 ResNet18 预期输入形状匹配的模拟输入数据
// (批大小, 通道数, 高, 宽)
val dummy_input = torch.randn(1, 3, 224, 224).to(device)
// 保存模型并返回大小的函数
def get_model_size(model, file_path="temp_model.pt"):
torch.save(model.state_dict(), file_path)
val size = os.path.getsize(file_path) / (1024 * 1024) // 大小,单位为 MB
os.remove(file_path)
return size
请务必使用 model.eval() 将模型设置为评估模式。这很重要,因为它会禁用 Dropout 等层,并使用运行统计数据对 BatchNorm 层进行归一化,这对于一致的推理和量化非常必要。
原始浮点模型的性能分析
在优化之前,我们先建立一个基准。我们将使用 torch.profiler.profile 来分析原始 FP32 模型的推理性能。性能分析工具会记录 CPU 和 GPU(如果可用)上不同操作的执行时间和内存消耗。
// 对 FP32 模型进行推理性能分析
println("正在分析 FP32 模型...")
with torch.profiler.profile(
activities=[
torch.profiler.ProfilerActivity.CPU,
torch.profiler.ProfilerActivity.CUDA, # 仅当 CUDA 可用时
],
record_shapes=true, // 可选:记录输入形状
profile_memory=true, // 可选:分析内存使用情况
with_stack=true // 可选:添加源代码上下文
) as prof:
with torch.profiler.record_function("model_inference"): // 标记此代码块
for _ <- range(10): // 运行多次推理以获得稳定测量结果
model_fp32(dummy_input)
// 打印按 self CPU 时间排序的性能分析结果
println("FP32 模型性能分析结果(按 self CPU 时间排序):")
println(prof.key_averages().table(sort_by="self_cpu_time_total", row_limit=10))
// 打印按 self CUDA 时间排序的性能分析结果(如果适用)
if device.type == 'cuda':
println("\nFP32 模型性能分析结果(按 self CUDA 时间排序):")
println(prof.key_averages().table(sort_by="self_cuda_time_total", row_limit=10))
// 获取基准推理时间(多次运行的平均值)
val start_time = time.time()
with torch.no_grad():
for _ <- range(50):
model_fp32(dummy_input)
val end_time = time.time()
val fp32_inference_time = (end_time - start_time) / 50
println(f"\nFP32 平均推理时间: {fp32_inference_time:.6f} 秒")
// 获取基准模型大小
val fp32_model_size = get_model_size(model_fp32)
println(f"FP32 模型大小: {fp32_model_size:.2f} MB")
查看性能分析工具的输出表格。查看 Name 列下消耗时间最多的操作(self_cpu_time_total 或 self_cuda_time_total)。对于 ResNet 等卷积网络,您通常会看到 aten::conv2d、aten::batch_norm、aten::relu 和 aten::addmm(用于线性层)占据了主要的执行时间。此分析印证了量化等优化工作可能带来最大益处的地方。
应用训练后静态量化(PTQ)
现在,我们将应用 PTQ,把 FP32 模型转换为量化的 INT8 版本。静态量化需要一个校准步骤,使用代表性数据来计算激活的量化参数(缩放因子和零点)。
注意: 对于 CPU 上的 PTQ,我们通常使用 ‘fbgemm’ 后端。对于 ARM CPU(移动设备上常见),通常优选 ‘qnnpack’。如果使用 CUDA,量化支持较为有限,通常依赖于特定的硬件功能(如 Tensor Cores)以及特定的后端或库(如 TensorRT)。为简化起见,本例侧重于使用 ‘fbgemm’ 进行 CPU 量化。
// --- 训练后静态量化 ---
println("\n正在开始训练后静态量化...")
// 创建模型的副本用于量化并移至 CPU
// 量化通常首先在 CPU 上执行
val quantized_model = copy.deepcopy(original_model)
quantized_model.eval()
quantized_model.cpu() // 将模型移至 CPU 以进行量化步骤
// 1. 模块融合:组合 Conv-BN-ReLU 序列以提升量化精度和性能
// 注意:融合列表可能需要根据具体的模型架构进行调整。
// 对于 ResNet,常见的融合包括 Conv-BN、Conv-BN-ReLU。
val modules_to_fuse = []
for (name, module) <- quantized_model.named_modules():
if module.isinstance(models.resnet.Bottleneck) || module.isinstance(models.resnet.BasicBlock):
// 寻找 (conv, bn, relu) 或 (conv, bn) 等序列
// 这是一个简化示例;实际实现可能需要更复杂的模式匹配。
val seq = []
for (child_name, child_module) <- module.named_children():
// 检查 Conv2d, BatchNorm2d, ReLU 模式
// 简单模式:conv -> bn -> relu 或 conv -> bn
if child_module.isinstance((torch.nn.Conv2d, torch.nn.BatchNorm2d, torch.nn.ReLU)):
seq.append(f"{name}.{child_name}")
if seq.length >= 2: // 找到了至少 conv-bn
// 检查最后两个是否为 Conv-BN
val is_conv_bn = module.get_submodule(seq(-2).split('.')(-1)).isinstance(torch.nn.Conv2d) && \
child_module.isinstance(torch.nn.BatchNorm2d)
if is_conv_bn:
// 如果 BN 后面跟着 ReLU,可选地添加 ReLU
val next_module_idx = module.named_children().index((child_name, child_module)) + 1
if next_module_idx < module.named_children().length:
val next_child_name, next_child_module = module.named_children()(next_module_idx)
if next_child_module.isinstance(torch.nn.ReLU):
modules_to_fuse.append(seq :+ f"{name}.{next_child_name}")
else:
modules_to_fuse.append(seq.copy())
else:
modules_to_fuse.append(seq.copy())
seq = [] // 找到匹配后重置序列
else: // 如果遇到不可融合的层,则中断序列
seq = []
// 也考虑顶层的 conv1, bn1, relu
if quantized_model.hasattr('conv1') && quantized_model.hasattr('bn1') && quantized_model.hasattr('relu'):
modules_to_fuse.append(Seq("conv1", "bn1", "relu"))
println(f"要融合的模块数量: ${modules_to_fuse.length}")
// 应用融合
if modules_to_fuse.nonEmpty:
quantized_model = torch.quantization.fuse_modules(quantized_model, modules_to_fuse, inplace=True)
println("模块融合完成。")
// 2. 指定量化配置
// 对 x86 CPU 使用 'fbgemm'。对 ARM CPU 使用 'qnnpack'。
quantized_model.qconfig = torch.quantization.get_default_qconfig('fbgemm')
print(f"量化配置设置为: {quantized_model.qconfig}")
// 3. 准备模型进行校准
// 插入观察器以收集激活统计信息
torch.quantization.prepare(quantized_model, inplace=True)
print("模型已准备好进行校准(观察器已插入)。")
// 4. 校准模型
// 在少量代表性数据集(校准数据)上运行推理
// 这里我们使用随机数据进行演示;在实际应用中,请使用验证集的一个子集。
println("正在运行校准...")
val calibration_data = [torch.randn(1, 3, 224, 224, dtype=torch.float32) for _ in range(100)] # 使用约 100 个样本
with torch.no_grad():
for input_data in calibration_data:
quantized_model(input_data)
println("校准完成。")
// 5. 将模型转换为量化版本
// 将模块替换为量化对应项并使用收集到的统计信息
val quantized_model = torch.quantization.convert(quantized_model, inplace=True)
println("模型已转换为量化版本 (INT8)。")
// 确保量化模型处于评估模式
quantized_model.eval()
让我们可视化简化的 PTQ 流程:
FP32 模型(预训练)融合模型(卷积-BN-ReLU)融合模块准备好的模型(已添加观察器)准备校准模型(已收集统计信息)校准INT8 模型(已量化)转换校准数据运行推理
该过程涉及融合兼容层,通过插入观察器准备模型,使用样本数据进行校准,最后转换为量化格式。
评估量化模型
现在,我们来评估 INT8 量化模型在 CPU 上的性能,并将其与原始 FP32 模型进行比较。我们将测量推理时间和模型大小。
// 对 CPU 上的 INT8 量化模型进行性能分析
println("\n正在分析 INT8 量化模型 (CPU)...")
// 确保模拟输入数据在 CPU 上用于量化模型
val dummy_input_cpu = dummy_input.cpu()
with torch.profiler.profile(
activities=[torch.profiler.ProfilerActivity.CPU], // 量化模型在此处在 CPU 上运行
record_shapes=true,
profile_memory=true,
with_stack=true
) as prof_quant:
with torch.profiler.record_function("quantized_model_inference"):
// 重要:确保已为量化操作设置后端
// 这通常是性能测量所必需的。
torch.backends.quantized.engine = 'fbgemm'
with torch.no_grad():
for _ in range(10):
quantized_model(dummy_input_cpu)
println("INT8 量化模型性能分析结果(按 self CPU 时间排序):")
println(prof_quant.key_averages().table(sort_by="self_cpu_time_total", row_limit=10))
println("INT8 量化模型性能分析结果(按 self CPU 时间排序):")
println(prof_quant.key_averages().table(sort_by="self_cpu_time_total", row_limit=10))
// 测量 INT8 推理时间
val start_time = time.time()
with torch.no_grad():
for _ in range(50):
quantized_model(dummy_input_cpu)
val end_time = time.time()
val int8_inference_time = (end_time - start_time) / 50
println(f"\nINT8 平均推理时间 (CPU): {int8_inference_time:.6f} 秒")
// 测量 INT8 模型大小
val int8_model_size = get_model_size(quantized_model)
println(f"INT8 模型大小: {int8_model_size:.2f} MB")
// --- 比较 ---
println("\n--- 性能比较 ---")
val speedup_factor = fp32_inference_time / int8_inference_time if device.type == 'cpu' else float('nan') // 仅直接比较 CPU 时间
val size_reduction = fp32_model_size / int8_model_size
println(f"FP32 推理所用设备: {device}")
println(f"FP32 平均推理时间: {fp32_inference_time:.6f} 秒")
println(f"INT8 平均推理时间 (CPU): {int8_inference_time:.6f} 秒")
if device.type == 'cpu':
println(f"CPU 推理加速比: {speedup_factor:.2f}x")
else:
println("CPU 推理加速比: 不适用 (FP32 在 GPU 上运行)")
println(f"\nFP32 模型大小: {fp32_model_size:.2f} MB")
println(f"INT8 模型大小: {int8_model_size:.2f} MB")
println(f"模型大小缩减: {size_reduction:.2f}x")
// 可选:可视化比较结果
import json
chart_data = {
"layout": {
"title": "模型性能比较",
"barmode": "group",
"xaxis": {"title": "指标"},
"yaxis": {"title": "数值"},
"font": {"family": "sans-serif"}
},
"data": [
{
"type": "bar",
"name": "推理时间 (秒)",
"x": ["FP32", "INT8 (CPU)"],
"y": [fp32_inference_time, int8_inference_time],
"marker": {"color": "#4dabf7"} # blue
},
{
"type": "bar",
"name": "模型大小 (MB)",
"x": ["FP32", "INT8 (CPU)"],
"y": [fp32_model_size, int8_model_size],
"marker": {"color": "#38d9a9"} # teal
}
]
}
// 根据数据类型正确格式化 y 轴以增加清晰度
chart_data["layout"]["yaxis"] = {"title": "时间 (秒) / 大小 (MB)"}
chart_data["layout"]["yaxis2"] = {
"title": "模型大小 (MB)",
"overlaying": "y",
"side": "right",
"showgrid": False,
}
// 将条形图分配给不同的轴
chart_data["data"][0]["yaxis"] = "y1"
chart_data["data"][1]["yaxis"] = "y2"
println("\n性能图表数据:")
println(f"```plotly\n{json.dumps(chart_data)}\n```")
原始 FP32 模型与 INT8 量化模型之间平均推理时间和模型大小的比较。请注意,只有当 FP32 模型也在 CPU 上运行时,直接的加速比比较才有意义。
讨论
本次实践练习展示了使用性能分析和训练后静态量化来优化 PyTorch 模型的标准流程。
- 性能分析: 我们使用
torch.profiler来确定原始 FP32 模型的性能特征。此步骤有助于理解计算时间的花费位置,并确认量化目标层(如卷积层)确实是重要的贡献者。 - 量化: 我们应用了 PTQ,其中包括融合模块、使用观察器准备模型、用样本数据校准,并将模型转换为 INT8。
- 评估: 将 INT8 模型与 FP32 基准进行比较通常会显示:
- 模型大小减小:INT8 权重和激活所需的存储空间显著减少(通常约 4 倍缩减)。
- 更快的 CPU 推理:INT8 操作可以在支持专用指令的 CPU 上更高效地执行,从而带来显著的加速。GPU 加速效果则在很大程度上取决于硬件支持和具体操作。
- 潜在的精度权衡:尽管 PTQ 旨在最大程度地减少精度损失,但仍可能发生一些性能下降。在您的特定任务和验证数据集上评估量化模型以确保其仍满足要求,这一点很重要。如果精度显著下降,那么之前讨论过的量化感知训练 (QAT) 等技术可能就会有必要了。
“这个动手示例为应用这些优化技术奠定了基础。请记住,具体步骤(如融合列表)和结果可能因模型架构、所选量化后端以及用于推理的硬件而异。尝试不同的配置并评估精度是部署场景中的重要后续步骤。”
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)