【NPU】Ascend Device Plugin — Part 2 gRPC 服务层 超深度源码分析之一
·
Ascend Device Plugin — gRPC 服务层 超深度源码分析
分析目录:
pkg/server/
涵盖文件:manager.go、plugin.go、server.go、types.go、pod_resource.go、dpu.go、manager_v2.go、plugin_v2.go、npu_base_v2.go、types_v2.go
一、模块定位
1.1 业务职责
pkg/server 是 Ascend Device Plugin 的 gRPC 服务核心层,承担以下关键职责:
| 职责 | 说明 |
|---|---|
| Kubernetes 设备插件 gRPC 服务 | 实现 Kubernetes Device Plugin v1beta1 接口(ListAndWatch、Allocate、GetDevicePluginOptions、PreStartContainer),向 kubelet 注册并上报设备状态 |
| 设备生命周期管理 | 管理华为昇腾 NPU(910A/910B/910A3/910A5/310/310P/310B)的发现、分类、健康监测、故障处理、热复位 |
| Pod 资源查询 | 通过 kubelet Pod Resources API 获取 Pod 已分配的设备信息 |
| 节点标签与注解管理 | 向 Kubernetes Node 写入标签(label)和注解(annotation),包括芯片名称、服务器类型、拓扑信息、SuperPod/Rack ID 等 |
| Volcano 调度集成 | 支持 Volcano 调度器场景下的设备分配映射(klt→real device 映射) |
| DPU 网卡健康监测 | 针对 A5(Ascend910A5)设备,周期性查询 DPU 网卡 operstate 并更新到 ConfigMap |
| 故障订阅与热复位 | 订阅 NPU/交换机故障事件,执行芯片热复位(hot reset),支持训练/推理场景的容错 |
| Rank Table 拓扑生成 | 为 A5 设备生成 HCCL Rank Table 所需的 LevelList 拓扑信息(URMA EID、RoCE 地址等) |
| vNPU 动态切分 | 支持 vNPU 的动态创建与销毁(仅 PresetVDevice=false 模式) |
| 软共享设备 | 支持 soft-share device 场景的配额管理和配置文件写入 |
1.2 在系统中的位置
1.3 文件职责总览
二、模块整体结构
2.1 核心类型与接口定义
2.2 核心方法清单
HwDevManager 核心方法
| 方法 | 文件 | 作用 |
|---|---|---|
NewHwDevManager |
manager.go | 构造函数,初始化设备管理器 |
setAscendManager |
manager.go | 根据设备类型设置运行模式和 manager |
setAllDeviceAndType |
manager.go | 初始化 K8s 客户端并获取所有 NPU 设备 |
UpdateNode / updateNode |
manager.go | 更新节点标签和注解 |
getNewNodeLabel |
manager.go | 计算需要写入的新节点标签 |
getNewNodeAnnotation |
manager.go | 计算需要写入的新节点注解(含 baseDevInfo、SuperPod 等) |
initPluginServer |
manager.go | 为每种设备类型创建 PluginServer 实例 |
ListenDevice |
manager.go | 主循环:订阅故障、周期性更新设备信息、通知 kubelet |
handleDeviceInfoUpdate |
manager.go | 周期设备信息更新核心逻辑 |
Serve |
manager.go | 管理 gRPC server 的启动与事件处理 |
notifyToK8s |
manager.go | 设备状态变化后通知 kubelet(通过 PluginServer.Notify) |
chipHotReset |
manager.go | 芯片热复位入口(推理卡) |
hotReset |
manager.go | 执行芯片复位并轮询等待启动完成 |
SubscribeFaultEvent |
manager.go | 订阅 NPU 和交换机故障事件 |
updatePodAnnotation |
manager.go | 更新 Pod 的 AscendReal 注解 |
useVolcanoNotify |
manager.go | 通过 Volcano ListAndWatch 上报设备信息 |
getSuperPodInfo |
manager.go | 获取 SuperPod 拓扑信息 |
ListenDpu |
dpu.go | 周期查询 DPU operstate |
updateDpuHealthy |
dpu.go | 更新 DPU 健康状态到 groupDevice |
getLevelList |
manager_v2.go | 生成 A5 Rank Table LevelList |
getROCEAddrList |
manager_v2.go | 获取 RoCE 地址列表 |
PluginServer 核心方法
| 方法 | 文件 | 作用 |
|---|---|---|
ListAndWatch |
plugin.go | gRPC 接口:向 kubelet 持续上报设备列表 |
Allocate |
plugin.go | gRPC 接口:为 Pod 分配设备并构造挂载响应 |
Notify |
plugin.go | 设备状态变更通知,触发 ListAndWatch 推送 |
responseToKubelet |
plugin.go | 构造 ListAndWatch 响应(含 volcano/非 volcano/动态切分三种模式) |
checkAllocateRequest |
plugin.go | 校验 Allocate 请求合法性 |
useVolcano / doWithVolcanoSchedule |
plugin.go | Volcano 调度场景下的设备分配 |
GetKltAndRealAllocateDev |
plugin.go | 获取 kubelet 分配设备与实际设备的映射 |
updateAllocMap |
plugin.go | 更新 klt→real device 映射表 |
DestroyNotUsedVNPU |
plugin.go | 销毁未使用的虚拟 NPU |
setNPUDeviceMount |
plugin.go | 设置设备挂载(ascend-docker 或原始模式) |
setHcclTopoFilePathEnv |
plugin_v2.go | 为 A5 设置 HCCL 拓扑文件路径环境变量 |
Start / Stop |
server.go | gRPC Server 启停与 kubelet 注册 |
serve |
server.go | 创建 gRPC server 并监听 unix socket |
register |
server.go | 向 kubelet 注册 Device Plugin |
PodResource 核心方法
| 方法 | 文件 | 作用 |
|---|---|---|
GetPodResource |
pod_resource.go | 调用 kubelet Pod Resources API 获取 Pod 设备分配 |
IsPodMoveComplete |
pod_resource.go | 检查故障设备上的 Pod 是否已完成迁移 |
getContainerResource |
pod_resource.go | 从容器资源中提取设备 ID 列表 |
NpuBase / ProductBase 核心方法
| 方法 | 文件 | 作用 |
|---|---|---|
getID |
npu_base_v2.go | 根据 rank level 返回网络实例 ID |
getRankLevelInfoKeyArr |
npu_base_v2.go | 获取各层级的网络类型数组(UB/UBG/UBoE/RoCE) |
getNetTypeAndFeIDListByRankLevel |
npu_base_v2.go | 根据层级获取网络类型和功能实体 ID 列表 |
getRandAddrByFuncEntityID |
npu_base_v2.go | 根据 phyID/feID/层级获取 RankAddr 列表 |
GetPortListByEid |
npu_base_v2.go | 根据 EID 获取端口列表(含缓存) |
getTopoFileInfo |
npu_base_v2.go | 读取并解析拓扑文件 |
SetUrmaDeviceInfoByHdm |
npu_base_v2.go | 通过 DCMI 获取 URMA 设备 EID 列表 |
2.3 内部调用关系
2.4 数据流入流出方式
三、核心业务逻辑深度解析
3.1 manager.go — HwDevManager 设备管理器核心
3.1.1 结构体定义
// HwDevManager 管理华为昇腾设备
type HwDevManager struct {
SwitchDevManager *deviceswitch.SwitchDevManager // 交换机设备管理器(A3 场景)
groupDevice map[string][]*common.NpuDevice // 按设备类型分组的设备列表
ServerMap map[string]InterfaceServer // 设备类型 → PluginServer 映射
allInfo common.NpuAllInfo // 所有 NPU 设备信息
manager device.DevManager // 设备管理器接口(910/310/310P)
RunMode string // 运行模式:Ascend310 / Ascend910 / Ascend310P
WorkMode string // 工作模式(910 的 SMP/AMP)
baseNPUInfo map[string]*common.NpuBaseInfo // NPU 基础信息缓存(用于注解比对)
dpuManager *dpucontrol.DpuFilter // DPU 过滤器(A5 场景)
ManagerLock sync.Mutex // 管理器锁
ContainerRuntime string // 容器运行时(docker/containerd)
}
// shareDevResourceQuota 软共享设备资源配额
type shareDevResourceQuota struct {
aicoreQuota int // AI Core 配额
hbmQuota int // HBM 内存配额
schedulingPolicy int // 调度策略
}
3.1.2 NewHwDevManager 构造函数
逐行解析:
func NewHwDevManager(devM devmanager.DeviceInterface) *HwDevManager {
var hdm HwDevManager
// 初始化 DPU 过滤器实例,用于 A5 场景的 DPU 信息管理
hdm.dpuManager = &dpucontrol.DpuFilter{}
// 步骤1:根据设备接口类型,设置对应的 manager(910/310/310P)
if err := hdm.setAscendManager(devM); err != nil {
hwlog.RunLog.Errorf("init hw dev manager failed, err: %v", err)
return nil // 初始化失败返回 nil,调用方需检查
}
// 步骤2:创建 K8s 客户端,获取所有 NPU 设备并分类
if err := hdm.setAllDeviceAndType(); err != nil {
hwlog.RunLog.Errorf("set all device and type failed, err: %v", err)
return nil
}
// 步骤3:初始化复位信息管理器,用于记录设备复位状态
device.InitResetInfoMgr(hdm.manager.GetKubeClient())
// 步骤4:检查产品类型是否支持(如 Atlas 300I Duo 不支持动态虚拟化)
if err := hdm.checkSupportedProductType(); err != nil {
hwlog.RunLog.Errorf("check supported product type failed, err: %v", err)
return nil
}
// 步骤5:获取并缓存 SuperPod 信息(SuperPod ID、Server Index、Rack ID)
hdm.setSuperPodInfo()
// 步骤6:更新 K8s 节点标签和注解(芯片名、服务器类型、拓扑信息等)
if err := hdm.UpdateNode(); err != nil {
hwlog.RunLog.Errorf("update node label failed, err: %v", err)
return nil
}
// 步骤7:如果设备类型从旧名称迁移到新名称(如 Ascend910 → huawei.com/npu),
// 则删除旧资源名,避免 kubelet 上残留旧资源
kubeClient := hdm.manager.GetKubeClient()
if kubeClient != nil {
deviceType := hdm.manager.GetDmgr().GetDevType()
if !customname.IsOldDeviceType(deviceType) {
hwlog.RunLog.Info("current device type changes to Huawei.com/npu, delete old resource name")
err := kubeClient.RemoveOldResource(api.HuaweiAscend910)
if err != nil {
hwlog.RunLog.Errorf("failed to delete old resource name: %v", err)
return nil
}
}
}
// 步骤8:为每种设备类型创建 PluginServer 实例
if err := hdm.initPluginServer(); err != nil {
hwlog.RunLog.Errorf("init plugin server failed, err: %v", err)
return nil
}
// 步骤9:获取容器运行时类型(docker/containerd)
if runtime, err := hdm.manager.GetKubeClient().GetContainerRuntime(); err == nil {
hdm.ContainerRuntime = runtime
}
return &hdm
}
3.1.3 setAscendManager — 设备类型路由
func (hdm *HwDevManager) setAscendManager(dmgr devmanager.DeviceInterface) error {
devType := dmgr.GetDevType()
// PresetVDevice=false 仅支持 310P 和 910B(动态切分模式)
if !common.ParamOption.PresetVDevice && devType != api.Ascend310P && devType != api.Ascend910B {
return fmt.Errorf("only 310p and 910b support to set presetVirtualDevice false")
}
common.ParamOption.RealCardType = devType // 缓存真实卡类型到全局参数
// 根据设备类型选择对应的 manager 实现
switch devType {
case api.Ascend310, api.Ascend310B:
hdm.RunMode = api.Ascend310
hdm.manager = device.NewHwAscend310Manager()
case api.Ascend910A, api.Ascend910B, api.Ascend910A3, api.Ascend910A5:
hdm.RunMode = api.Ascend910
hdm.manager = device.NewHwAscend910Manager()
hdm.WorkMode = dmgr.GetNpuWorkMode() // 910 需要额外的工作模式
case api.Ascend310P:
hdm.RunMode = api.Ascend310P
hdm.manager = device.NewHwAscend310PManager()
default:
return fmt.Errorf("an unsupported device type")
}
hdm.manager.SetDmgr(dmgr) // 注入设备管理接口
// 获取所有产品类型并缓存
productTypes, err := hdm.manager.GetDmgr().GetAllProductType()
if err != nil {
return err
}
common.ParamOption.ProductTypes = productTypes
// 检查 310P 混插模式是否合法
if err = common.CheckCardUsageMode(common.ParamOption.Use310PMixedInsert, productTypes); err != nil {
return err
}
// 非边缘场景需要获取芯片 AI Core 数量
if common.ParamOption.BuildScene != common.EdgeScene {
aiCoreCount, err := hdm.manager.GetChipAiCoreCount()
if err != nil {
return err
}
common.ParamOption.AiCoreCount = aiCoreCount
}
return nil
}
3.1.4 ListenDevice — 主循环
这是整个 device plugin 的 核心运行循环,启动后持续运行直到收到停止信号。
逐行解析关键部分:
func (hdm *HwDevManager) ListenDevice(ctx context.Context) {
hwlog.RunLog.Info("starting the listen device")
// 1. 订阅故障事件(NPU + 交换机)
hdm.subscribeFaultEvent()
// 2. 如果是 A3 设备且启用了交换机故障检测,启动定期查询 goroutine
if common.ParamOption.RealCardType == api.Ascend910A3 && common.ParamOption.EnableSwitchFault {
go hdm.SwitchDevManager.GetSwitchFaultCodeByInterval(ctx,
time.Second*common.GetSwitchFaultCodeInterval)
}
// 3. 加载故障码和设备信息 ConfigMap(阻塞,完成后启动后台轮询)
hdm.loadFaultCodeAndDeviceInfoCm(ctx)
// 4. 启动 gRPC Server 管理 goroutine(处理 socket 文件监听、重启等)
go hdm.Serve(ctx)
// 5. 如果需要检查缓存的 Pod,启动 Pod Informer 检查器
if common.ParamOption.CheckCachedPods {
go hdm.manager.GetKubeClient().PodInformerInspector(ctx)
}
// 6. 启动节点注解定期更新 goroutine(每 60 秒检查一次)
go hdm.updateNodeAnnotations(ctx)
// 7. 启动设备故障写入 K8s Event goroutine
go hdm.manager.WriteFaultToEvent(ctx)
// 8. 主循环:周期性设备信息更新 + 触发器检查
initTime := time.Now()
ticker := time.NewTicker(time.Duration(common.ParamOption.ListAndWatchPeriod) * time.Second)
defer ticker.Stop()
triggerTicker := time.NewTicker(time.Second) // 每秒检查触发器
defer triggerTicker.Stop()
for {
select {
case _, ok := <-ctx.Done():
// 收到上下文取消信号,退出主循环
hwlog.RunLog.Info("listen device stop")
return
case <-triggerTicker.C:\n // 每秒检查是否有更新触发信号(如故障事件触发的即时更新)\n hdm.parseTriggers(ctx, initTime)\n case <-ticker.C:\n // 周期性设备信息更新(默认 5 秒)\n hwlog.RunLog.Debug("Periodic device info update")
hdm.handleDeviceInfoUpdate(ctx, &initTime)
}
}
}
3.1.5 handleDeviceInfoUpdate — 周期设备信息更新
逐行解析:
func (hdm *HwDevManager) handleDeviceInfoUpdate(ctx context.Context, initTime *time.Time) {
// 加全局锁,防止并发更新设备信息
common.LockAllDeviceInfo()
defer common.UnlockAllDeviceInfo()
// 1. 更新所有设备信息(动态切分场景需先销毁未用 vNPU,重新获取设备列表)
if err := hdm.updateAllInfo(); err != nil {
hwlog.RunLog.Error(err)
return
}
// 2. 补充订阅方式无法上报的故障码(如断卡、丢芯片等)
hdm.mendSubscribeFaultEvents()
// 3. 更新 Pod 注解(写入 AscendReal 真实分配设备信息)
if err := hdm.updatePodAnnotation(); err != nil {
hwlog.RunLog.Error(err)
}
// 4. 标记设备是否被 Pod 使用(从 kubelet 获取已分配设备列表)
hdm.updateDeviceUsedInfo(hdm.groupDevice)
// 5. 通知 kubelet 设备状态变化(触发 ListAndWatch 推送)
hdm.notifyToK8s(ctx, initTime)
// 6. 检查节点复位信息:如果 annotation 中有复位失败设备但设备已恢复,清除注解
hdm.checkNodeResetInfo()
// 7. Volcano ListAndWatch 上报(设备健康状态、内存信息)
hdm.useVolcanoNotify()
// 8. 推理卡热复位检查(仅 HotResetInfer 模式)
hdm.chipHotReset()
// 9. 删除已恢复的故障记录和频率故障记录
common.DelOnceRecoverFault(hdm.groupDevice)
common.DelOnceFrequencyFault()
// 10. 标记同步完成,允许其他操作读取最新设备信息
common.Synchronize = true
}
3.1.6 notifyToK8s — 通知 kubelet 设备状态
func (hdm *HwDevManager) notifyToK8s(ctx context.Context, initTime *time.Time) {
// 检查是否支持训练容错
hdm.isSupportGraceTolerance()
// 深拷贝旧设备状态(用于变更检测)
oldGroupDevice := deepCopyGroupDevice(hdm.groupDevice)
// 更新设备健康状态(根据故障码判断 healthy/unhealthy)
hdm.manager.UpdateHealth(hdm.groupDevice, hdm.allInfo.AICoreDevs, hdm.RunMode)
// A5 设备:更新 DPU 健康状态
if hdm.manager.GetDmgr().GetDevType() == api.Ascend910A5 {
hdm.updateDpuHealthy(hdm.groupDevice)
}
// 训练容错处理(热复位场景下,正在复位的设备标记为 healthy)
hdm.graceTolerance(ctx, hdm.groupDevice)
// 检测设备状态是否变化
isDevStateChange := hdm.manager.GetChange(hdm.groupDevice, oldGroupDevice)
// 对每种设备类型,如果状态变化或需要重发,通知 PluginServer
for devType, isChanged := range isDevStateChange {
server := hdm.ServerMap[devType]
if server == nil {
continue
}
// 如果状态未变化,且满足以下条件之一则跳过:
// - 启动 1 分钟内且上次发送成功,且启动 1 小时内
if !isChanged &&
(time.Now().Sub(*initTime) < time.Minute || server.LastSendSuccess()) &&
time.Now().Sub(*initTime) < time.Hour {
continue
}
*initTime = time.Now()
// 非预置虚拟设备:使用 AICore 设备列表通知
if !common.ParamOption.PresetVDevice {
hdm.pluginNotify(hdm.allInfo.AICoreDevs, common.AiCoreResourceName)
return
}
// 预置虚拟设备:按设备类型分别通知
hdm.pluginNotify(hdm.groupDevice[devType], devType)
}
}
3.1.7 hotReset — 芯片热复位
func (hdm *HwDevManager) hotReset(device *common.NpuDevice, devices []*common.NpuDevice) {
hwlog.RunLog.Infof("will start to reset device %s", device.DeviceName)
// 标记设备正在复位,防止重复复位
hdm.manager.SetCardsInResetting(device.LogicID, true)
var isResetExec = false
successResetDevList := sets.NewInt32()
// 轮询等待:每秒检查一次,最多等待 1 分钟
if err := wait.PollImmediate(time.Second, time.Minute, func() (bool, error) {
// 执行复位操作(只执行一次,通过 isResetExec 标志控制)
if err := hdm.execResetChip(device.LogicID, &isResetExec); err != nil {
return false, err
}
// 检查所有关联设备的启动状态
for _, dev := range devices {
if successResetDevList.Has(dev.LogicID) {
continue // 已确认启动完成的设备跳过
}
bootState, err := hdm.manager.GetDmgr().GetDeviceBootStatus(dev.LogicID)
if err != nil {
return false, err
}
if bootState != common.BootStartFinish {
// 设备尚未启动完成,继续等待
return false, nil
}
successResetDevList.Insert(dev.LogicID)
}
// 所有设备都启动完成
common.SetDeviceInit(device.LogicID)
return true, nil
}); err != nil {
// 超时或出错
hwlog.RunLog.Warnf("hot reset failed, timeout or err: %v", err)
hdm.manager.SetCardsInResetting(device.LogicID, false)
hdm.manager.SetResetFailedTimes(device.LogicID, hdm.manager.GetResetFailedTimes(device.LogicID)+1)
return
}
// 复位成功:清除复位失败计数,清除复位中标记
hdm.manager.SetResetFailedTimes(device.LogicID, 0)
hdm.manager.SetCardsInResetting(device.LogicID, false)
hwlog.RunLog.Info("hot reset success")
}
3.1.8 Serve — gRPC Server 管理
func (hdm *HwDevManager) Serve(ctx context.Context) {
hwlog.RunLog.Info("Serve start")
// 创建 socket 文件监听器(监听 /var/lib/kubelet/device-plugin/ 目录)
watcher, err := common.NewFileWatch()
if err != nil {
hwlog.RunLog.Error("createSocketWatcher error")
return
}
defer func() {
if watcher != nil {
if err := watcher.FileWatcher.Close(); err != nil {
hwlog.RunLog.Errorf("close file watcher, err: %v", err)
}
}
}()
// 创建重启信号监听(SIGHUP)
restartSignal := common.NewSignWatcher(syscall.SIGHUP)
for {
// 启动所有标记需要重启的 PluginServer
allSuccess := hdm.startAllServer(watcher)
// 等待事件(上下文取消、重启信号、socket 文件变更)
if hdm.handleEvents(ctx, restartSignal, watcher) {
break // 收到停止信号,退出循环
}
// 如果启动失败,等待一段时间后重试
if !allSuccess {
time.Sleep(common.SleepTime * time.Second)
}
}
}
handleEvents 处理三种事件:
- ctx.Done():收到停止信号,返回 true 退出
- restartSignal (SIGHUP):设置所有 server 重启标志
- fsnotify 事件:socket 文件被删除时打印告警;kubelet.sock 创建时记录日志
3.1.9 updateNode — 节点标签与注解更新
getNewNodeLabel 关键逻辑:
updateChipNameToNode():获取芯片名称写入标签- 服务器类型标签:
Ascend910-32或ascend-32(新命名) - HBM 内存标签:
NPUChipMemoryLabel(如32G) - A300IA2 推理卡标签:当 910B + 推理模式 + 特定 BoardId 时添加
- 拓扑标签:A3/A5 添加 SuperPodId,A5 额外添加 RackId
getNewNodeAnnotation 关键逻辑:
cardType:板卡类型(如 A5300I)BaseDevInfoAnno:JSON 序列化的 NPU 基础信息(设备名→IP/DeviceID/SuperDeviceID/LevelList)SuperPodIDKey/serverIndexKey/serverTypeKey/RackIDKey:SuperPod 拓扑信息
3.1.10 故障码 ConfigMap 管理
3.1.11 chipHotReset — 推理卡热复位
鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。
更多推荐

所有评论(0)