3.2 ascend910.go — HwAscend910Manager 深度解析

3.2.1 结构体与构造函数
type HwAscend910Manager struct {
    AscendTools                              // 嵌入基类
    dpu             common.DpuInfo           // DPU信息(总线类型/DPU列表/NPU→DPU映射)
    hotResetManager HotResetManager          // 热重置管理器接口
}
func NewHwAscend910Manager() *HwAscend910Manager {
    return &HwAscend910Manager{
        AscendTools: AscendTools{
            name:                      getAscend910Name(),  // A5→"npu", 其他→"Ascend910"
            unHealthyKey:              common.GetAscend910Key(api.CmCardUnhealthySuffix),
            devCount:                  common.MaxDevicesNum,
            cardInResetMap:            make(map[int32]bool, common.GeneralMapSize),
            resetFailedTimesMap:       make(map[int32]int, common.GeneralMapSize),
            lastUsedChipsContainerMap: make(map[string]sets.String),
        },
    }
}

设计意图

  • hotResetManager 通过接口而非具体类型引用,便于测试mock和未来扩展
  • dpu 字段独立于AscendTools,因为DPU是910A5特有的概念
  • 构造函数只初始化基本字段,hotResetManagerGraceTolerance首次调用时通过sync.Once延迟初始化
3.2.2 核心流程:GetNPUs — 设备发现

是(物理设备)

否(有虚拟设备)

GetNPUs()

dmgr.GetDeviceList()
获取设备数量和逻辑ID列表

devNum >
devCount?

错误: 设备数超限

遍历每个逻辑ID

getDavinCiDev(logicID)
获取物理ID/卡ID/设备ID/IP

getVirtualDevice(logicID)
查询虚拟设备信息

VDevNum >
MaxVirtualDeviceNum?

错误: 虚拟设备数超限

PresetVDevice?
预置虚拟设备?

FakeAiCoreDevice
创建AI Core虚拟设备

跳过

VDevNum == 0?

assemblePhyDevices
组装物理设备
名称: Ascend910-{phyID}

assembleVirtualDevices
组装虚拟设备
名称: {type}-{vDevType}-{vDevID}-{phyID}

removeDuplicate
设备类型去重

返回NpuAllInfo
{AllDevs, AICoreDevs, AllDevTypes}

虚拟设备命名规则

物理设备:  Ascend910-{phyID}                    例如: Ascend910-0
虚拟设备:  {name}-{vDevType}-{vDevID}-{phyID}   例如: Ascend910-4c-1-0
共享设备:  {name}-{phyID}-{index}               例如: Ascend910-0-0
3.2.3 核心流程:GraceTolerance — 优雅容错主入口

HotResetTrainOnLine
(在线容错)

其他(离线容错)

setTaskDevInfoCache 详细

GetActivePodListCache
获取活跃Pod列表

遍历Pod

获取华为Ascend910注解

convertPhysicIdToLogicId
物理ID→逻辑ID

isReSchedulingScene
重调度场景检查

GetTaskNameByPod
获取任务名

获取RankIndex注解

GenerateTaskDevFaultInfoList
生成任务设备故障信息

UpdateFaultDev2PodMap
更新故障设备→Pod映射

handleUpdateCaches
批量更新缓存

updateHotResetCache 详细

updateUpgradeErrorInfo
更新隔离设备列表

UpdateGlobalDevFaultInfoCache
更新全局设备故障缓存

setTaskDevInfoCache
设置任务设备信息缓存

GraceTolerance(ctx, classifyDevs)

hotResetManagerInitOnce.Do
首次调用初始化HotResetManager

hotResetManager
== nil?

日志错误, return

hotResetManager.SyncResetCM
启动CM/Pod Informer同步

GraceToleranceOn?

return

updateHotResetCache
更新热重置缓存

更新成功?

日志错误, return

HotReset模式?

processAllTask
处理所有任务容错

处理成功?

日志错误

hotResetHandler
(goroutine)
无任务热重置

filterDevStatus
过滤重置中设备状态

setAllDevUnhealthyOnRing
设置Ring上重置设备不健康

updateHotResetCache 逐行解析

func (hnm *HwAscend910Manager) updateHotResetCache(classifyDevs map[string][]*common.NpuDevice) error {
    deviceList, ok := classifyDevs[hnm.name]
    if !ok {
        return fmt.Errorf("ascend npu device list not found")
    }
    
    // 1. 更新已隔离设备列表
    //    遍历isolateDevList,如果设备已恢复健康则从隔离列表中移除
    if err := hnm.updateUpgradeErrorInfo(classifyDevs); err != nil {
        return err
    }
    
    // 2. 更新全局设备故障信息缓存
    //    为每个设备创建DevFaultInfo,设置Policy(隔离/正常/各故障级别)
    if err := hnm.hotResetManager.UpdateGlobalDevFaultInfoCache(deviceList, isolateDevList); err != nil {
        return err
    }
    
    // 3. 设置任务设备信息缓存
    //    扫描所有活跃Pod,构建任务→设备列表/故障信息/Pod的映射
    if err := hnm.setTaskDevInfoCache(); err != nil {
        return err
    }
    
    return nil
}
3.2.4 核心流程:processAllTask — 任务级容错处理

FreeResetError或更低

RestartRequestError(L2)

RestartError(L3)

ResetError(L5)

任务在重置中

否 + 有故障设备

否 + 无故障

任务不在重置中

L2

L3

L5

processAllTask(classifyDevs)

GetAllTaskDevFaultInfoList
获取所有任务故障信息

遍历每个任务

GetTaskProcessPolicy
获取任务最高级别策略

policyLevelHandle
策略级别检查

跳过此任务

日志: 开始处理L2

日志: 开始处理L3

isHotResetOn = true
日志: 开始处理L5

isolateSceneHandle
隔离场景检查

当前节点
在重置?

tryWriteIsolationInfo
(goroutine)
写隔离信息到CM

跳过此任务

preProcess
预处理

GetTaskDevFaultInfoList

GetTaskPod

GetTaskResetInfo
(policy, UnrecoveredStatus)

WriteResetInfoDataIntoCM
写ResetInfo CM

WriteFaultInfoDataIntoCM
写FaultInfo CM(可选)

SetTaskInReset

SetAllDevInReset

runProcessTask
根据策略级别启动处理

go restartRequestProcess

go restartProcess

go resetProcess

3.2.5 核心流程:L2故障处理 (restartRequestProcess)

超时
(WaitErrorCodeCleanTime秒)

restartRequestProcess
(goroutine)

defer: postProcess
清理重置状态

GetTaskDevFaultInfoList

RecordFaultInfoList
记录故障信息

DeepCopyDevFaultInfoList
深拷贝(用于CM写入)

checkDevErrorCode
等待错误码清除

PollImmediate
每1秒检查

refreshDevFaultInfo
刷新故障信息

GetDevListByPolicyLevel
(RestartRequestErrorLevel)

故障列表
为空?

L2自愈成功!

sleep 1秒

needUpgrade = true

handleSucceedRestartRequest
SetDeviceInit + 更新CM为Recovered + UnSetTaskInReset

needUpgrade?

return

upgradeRestartRequestProcess
L2升级为重置

updateResetCMStatusWithoutWait
(ResetError, Unrecovered)

waitForAllFaultyDeviceProcessesToZero
等进程归零

resetDeviceOnce
执行设备重置

仍有L2故障?

L2升级重置后恢复

updateResetCMStatus
(IsolateError, RecoverFailed)

return error

完成

3.2.6 核心流程:L3故障处理 (restartProcess)

失败

成功

失败

成功

restartProcess
(goroutine)

defer: postProcess

GetTaskDevFaultInfoList

updateResetCMStatus
(IsolateError, RecoverFailed)

RecordFaultInfoList

DeepCopyDevFaultInfoList

waitForAllFaultyDeviceProcessesToZero
等故障设备进程归零

日志错误, return

sleep WaitFaultSelfHealingTime
等待故障自愈

refreshDevFaultInfo

upgradeRestartProcess
检查是否需要升级为重置

有L3故障?

L3自愈成功!

updateResetCMStatusWithoutWait
(ResetError, Unrecovered)

resetDeviceOnce

仍有L3故障?

重置后L3恢复

updateResetCMStatus
(IsolateError, RecoverFailed)

updateResetCMStatus
(RestartError, Recovered)

UnSetTaskInReset

完成

3.2.7 核心流程:L5故障处理 (resetProcess)

失败

成功

失败

成功

resetProcess
(goroutine)

defer: isHotResetOn=false
+ postProcess

GetTaskDevFaultInfoList

RecordFaultInfoList

DeepCopyDevFaultInfoList

waitForAllFaultyDeviceProcessesToZero

updateResetCMStatus
(IsolateError, RecoverFailed)

resetDeviceOnce
直接执行设备重置

updateResetCMStatus
(IsolateError, RecoverFailed)

upgradeResetProcess
检查是否需要升级为隔离

有需重置
设备?

updateResetCMStatus
(ResetError, Recovered)

updateResetCMStatus
(IsolateError, RecoverFailed)

UnSetTaskInReset

完成(隔离)

完成(恢复)

3.2.8 故障升级链

超时未清除

重置后仍有故障

重启后仍有故障

重置后仍有故障

重置后仍有故障

L2: RestartRequestError
故障自愈
(等待错误码清除)

升级为ResetError
等待进程归零→重置设备

IsolateError
隔离

L3: RestartError
设备重启
(等进程归零→重启)

升级为ResetError
执行设备重置

IsolateError
隔离

L5: ResetError
强制重置
(等进程归零→重置)

IsolateError
隔离

3.2.9 核心流程:设备重置执行 (execResetDevice)

失败

成功

失败

成功

execResetDevice(devMap)

GetNeedResetDevMap
获取需重置设备映射

A3设备?

getNeedResetDevMapForA3

使用原映射

Ascend910ResetGoroutine.Store
记录重置协程PID

遍历devMap

canResetDevice
检查可重置(非忙/未超限)

跳过

isShouldCheckNet
是否需检查网络(分布式?)

tryResetDevice
尝试带内重置

加入失败列表

isRingResetComplete
等待Ring重置完成

FreeBusyDev
释放忙状态

加入成功列表

全部成功?

updateResetInfo
更新重置信息文件

A3设备?

execOutBandReset
带外重置失败设备

scanDeviceForThirdParty
(goroutine)
延迟第三方扫描

Ascend910ResetGoroutine.Delete

3.2.10 设备重置详细流程

isRingResetComplete (Ring完成检查)

!= BootStartFinish

== BootStartFinish

计算Ring范围
startID ~ startID+resetDevNumOnce

遍历Ring上每个设备

waitDeviceResetComplete
等待单个设备完成

PollImmediate 每秒检查

GetDeviceBootStatus
获取启动状态

shouldCheckNet?

设备完成

isNetResetCompleted
网络健康码=0或6?

resetDeviceOutBand (带外重置)

GetOutBandChannelState
检查带外通道

PreResetSoc
预重置SoC

SetDeviceResetOutBand
带外重置

sleep beforeRescanDelay
3秒等待

RescanSoc
重新扫描SoC

成功

tryResetDeviceOffline (离线重置)

不存在

存在

成功

失败

AddResetCnt

AddBusyDev

循环 ResetRetryTimes 次

checkFaultIsExist
检查故障是否仍存在

停止重置, 返回成功
(故障已自愈)

dmgr.SetDeviceReset

返回成功

sleep 指数退避

tryResetDevice (带内重置)

成功

失败

AddResetCnt
增加重置计数

AddBusyDev
标记设备忙

循环 ResetRetryTimes 次

dmgr.SetDeviceReset
DCMI重置接口

返回成功

sleep (i+1)*ResetInterVal
指数退避

3.2.11 核心流程:DoWithVolcanoListAndWatch

DoWithVolcanoListAndWatch
(classifyDevs, chipMemory)

getDevStatesDevSet
(AscendTools基类方法)

UpdateNodeDeviceInfo
(devStatusSet, dpu, updateDeviceInfo)

调用 hnm.updateDeviceInfo
(具体更新逻辑)

getHealthAndRecoverDev
获取健康设备和恢复设备

AutoStowingDevs?
自动收纳?

恢复标签置空
健康设备=全量健康

计算恢复集
= 旧不健康 - 新不健康

getNewNetworkRecoverDev
获取网络恢复设备

设置CM各Key的值

Ascend910: 健康设备列表

Ascend910-Recovering: 重置中设备

Ascend910-Unhealthy: 不健康设备

Ascend910-NetworkUnhealthy: 网络不健康

Ascend910-DPUUnhealthy: DPU不健康

Ascend910-FaultList: 故障详情JSON

isNeedBlockAllDevice?
(推理+A800IA2+有故障)

清空健康设备
所有设备标记为Recovering
阻止新调度

AutoStowingDevs?

返回(不更新NodeLabel)

update910NodeLabel
更新节点恢复/网络恢复标签

3.2.12 核心流程:hotResetHandler — 无任务热重置

hotResetHandler(classifyDevs)
(goroutine)

isHotResetOn?

跳过

获取910设备列表

遍历设备

getDevFaultInfo
从缓存获取故障信息

故障信息
为nil?

resetTimeMap.Delete
清除重置时间

getResetIndex
获取Ring索引

Ring已
有设备待重置?

跳过(同Ring只重置一次)

canBeReset
检查Ring上所有芯片空闲

isFaultNeedRestart
故障是否需重启

标记Ring
加入待重置列表

遍历待重置设备

startUpHotReset
启动热重置

handleResetProcess

hotResetTryOutBand
尝试带外重置

isHotResetOn = false

完成

isFaultNeedRestart 逐行解析

func (hnm *HwAscend910Manager) isFaultNeedRestart(devFaultInfo *common.DevFaultInfo) bool {
    // 910B系列:FreeResetError或ResetError直接需要重启
    if common.ParamOption.RealCardType == api.Ascend910B &&
        (devFaultInfo.Policy == common.FreeResetError || 
         devFaultInfo.Policy == common.ResetError) {
        return true
    }
    
    // 其他系列:检查故障持续时间
    resetTime := getResetTime(devFaultInfo.LogicId)
    if resetTime == 0 {
        // 首次检测到故障,记录开始时间,不立即重置
        resetTimeMap.Store(devFaultInfo.LogicId, time.Now().Unix())
        return false
    }
    
    // 故障存在超过 ResetFaultToleranceTimeInterval(60秒) → 需要重置
    if time.Now().Unix()-resetTime > common.ResetFaultToleranceTimeInterval {
        resetTimeMap.Delete(devFaultInfo.LogicId)
        return true
    }
    return false
}
3.2.13 A3卡关联重置机制

filterDevStatusForA3

遍历设备

设备在resetDev中?
且不健康?
且不应隔离?

设为Healthy
(重置中不报不健康)

GetAssociatedLogicIDs

关联设备网络也设为Healthy

跳过

setUnhealthyForA3

获取inResetDev

GetAssociatedLogicIDs

遍历关联LogicID

设置Health=Unhealthy
NetworkHealth=Unhealthy
Status=NPUResettingStatus

GetAssociatedLogicIDs

输入: logicID, cardID, deviceID

GetBrotherCardID
获取兄弟卡ID

GetDeviceLogicID(brotherCard, 0)

GetDeviceLogicID(brotherCard, 1)

otherDeviceId = (deviceID+1)%2

GetDeviceLogicID(cardID, otherDeviceId)

返回[logicID, ringDevLogic, logicID0, logicID1]

A3卡拓扑

物理卡A
cardID=X

die0
deviceID=0
logicID=L0

die1
deviceID=1
logicID=L1

物理卡B(兄弟卡)
cardID=Y
(GetBrotherCardID)

die0
deviceID=0
logicID=L2

die1
deviceID=1
logicID=L3

Ring(4个设备)

设计意图:A3卡每张物理卡包含2个die(deviceID=0和1),两张物理卡组成一个Ring。重置时必须同时重置Ring上所有4个设备。GetAssociatedLogicIDs 计算出完整的关联设备列表。

3.2.14 核心流程:filterDevStatus — 重置中设备状态过滤

filterDevStatus(classifyDevs)

获取设备列表

A3设备?

filterDevStatusForA3

普通Ring过滤

遍历设备

在resetDev中?
且不健康?
且不应隔离?

设为Healthy

跳过

GetAssociatedLogicIDs

关联设备网络设为Healthy

GetDevListInReset
获取重置中设备

GetResetDevNumOnce
Ring大小

遍历设备

在resetDev中?
且不健康?
且不应隔离?

Health = Healthy

跳过

计算Ring索引

遍历Ring上设备

NetworkHealth = Healthy

设计意图:当设备正在热重置时,其健康状态会被过滤为Healthy。这是因为:

  1. 设备正在重置过程中,故障可能已清除但尚未完全恢复
  2. 避免Volcano调度器将重置中的设备标记为不健康,导致任务被驱逐
  3. 重置完成后,设备状态由setAllDevUnhealthyOnRing统一管理
3.2.15 核心流程:waitForAllFaultyDeviceProcessesToZero
func (hnm *HwAscend910Manager) waitForAllFaultyDeviceProcessesToZero(taskName string,
    devFaultInfoList []*common.TaskDevInfo) error {
    // 1. 获取需要重置的设备LogicID映射
    faultDeviceLogicIdMap, err := hnm.getNeedResetDeviceLogicIdMap(devFaultInfoList)
    
    // 2. 创建超时定时器
    timer := time.NewTimer(common.WaitProcessesToZeroTime * time.Second)
    defer timer.Stop()
    
    timeCount := 0
    for {
        select {
        case <-timer.C:\n            // 超时:最后一次检查,如果仍非零→升级为隔离\n            if hnm.canContinueGraceProcess(faultDeviceLogicIdMap, taskName, true) {\n                return nil  // 最后一次检查通过
            }
            // 超时且进程未归零→更新CM为隔离状态
            hnm.updateResetCMStatusToIsolate(taskName, devFaultInfoList)
            return fmt.Errorf("check processes timeout")
            
        default:
            // 非阻塞检查:每PollingInterval秒检查一次
            if hnm.canContinueGraceProcess(faultDeviceLogicIdMap, taskName, false) {
                return nil  // 所有进程归零
            }
            time.Sleep(common.PollingInterval * time.Second)
            timeCount++
        }
    }
}

canContinueGraceProcesscheckNumberOfAllProcessIsZero

func (hnm *HwAscend910Manager) checkNumberOfAllProcessIsZero(
    faultDeviceLogicIdMap map[int32]int32) (bool, error) {
    if len(faultDeviceLogicIdMap) == 0 {
        return true, nil
    }
    isNumberOfAllProcessZero := true
    for faultDeviceLogicId, num := range faultDeviceLogicIdMap {
        if num == 0 {
            continue  // 已知为0的设备跳过查询(优化:减少DCMI调用)
        }
        // 通过DCMI获取设备进程信息
        devProcessInfo, err := hnm.dmgr.GetDevProcessInfo(faultDeviceLogicId)
        if err != nil || devProcessInfo == nil {
            return false, err
        }
        // 更新映射中的进程数(下次循环可跳过已为0的设备)
        faultDeviceLogicIdMap[faultDeviceLogicId] = devProcessInfo.ProcNum
        isNumberOfAllProcessZero = devProcessInfo.ProcNum == 0
        if !isNumberOfAllProcessZero {
            return false, nil  // 只要有一个设备进程非0,就返回false
        }
    }
    return isNumberOfAllProcessZero, nil
}
3.2.16 核心流程:updateDeviceInfo — 910设备信息CM更新
func (hnm *HwAscend910Manager) updateDeviceInfo(oldDevInfo, newDevInfo map[string]string,
    devStatusSet common.DevStatusSet) error {
    // 1. 获取健康设备和需要恢复的设备
    nodeFmtDevRecover, nodeFmtDevNetRecover := sets.String{}, sets.String{}
    newDevRecoverLabel, newAscend910 := hnm.getHealthAndRecoverDev(devStatusSet, 
        nodeFmtDevRecover,
        common.ConvertDevListToSets(oldDevInfo[common.GetAscend910Key(api.CmCardUnhealthySuffix)],
            common.CommaSepDev))
    
    // 2. 获取网络恢复设备
    newNetRecoverSets, newNetUHDevSets := hnm.getNewNetworkRecoverDev(
        devStatusSet.NetUnHealthyDevice,
        common.ConvertDevListToSets(oldDevInfo[common.GetAscend910Key(api.CmCardNetworkUnhealthySuffix)],
            common.CommaSepDev),
        nodeFmtDevNetRecover)
    
    // 3. 写入各Key的设备列表
    newDevInfo[common.GetAscend910Key("")] = newAscend910                    // 健康设备
    newDevInfo[common.GetAscend910Key(api.CmRecoveringSuffix)] = 
        common.ToString(devStatusSet.RecoveringDevices, common.CommaSepDev)   // 重置中设备
    
    // 4. 推理场景:A800IA2 with HCCS且有故障→阻塞所有设备调度
    if common.ParamOption.HotReset == common.HotResetInfer &&
        hnm.GetResetFailedTimes(common.FirstDevice) <= common.MaxResetTimes &&
        hnm.isNeedBlockAllDevice(devStatusSet.DeviceFault) {
        newDevInfo[common.GetAscend910Key("")] = ""                           // 清空健康设备
        newDevInfo[common.GetAscend910Key(api.CmRecoveringSuffix)] = 
            common.ToString(devStatusSet.AllDevices, common.CommaSepDev)      // 所有设备标记为Recovering
    }
    
    // 5. 写入不健康/网络不健康/DPU不健康设备列表
    newDevInfo[common.GetAscend910Key(api.CmCardUnhealthySuffix)] = 
        common.ToString(devStatusSet.UnHealthyDevice, common.CommaSepDev)
    newDevInfo[common.GetAscend910Key(api.CmCardNetworkUnhealthySuffix)] = 
        common.ToString(newNetUHDevSets, common.CommaSepDev)
    newDevInfo[common.GetAscend910Key(api.CmCardDPUUnhealthySuffix)] = 
        common.ToString(devStatusSet.DpuUnHealthyDevice, common.CommaSepDev)
    
    // 6. 写入故障详情JSON
    data := common.MarshalData(devStatusSet.DeviceFault)
    newDevInfo[common.GetAscend910Key(api.CmFaultListSuffix)] = string(data)
    
    // 7. 非自动收纳模式:更新Node Label
    if !common.ParamOption.AutoStowingDevs {
        curNode, err := hnm.getRecoverLabelFromNodeSets(&nodeFmtDevRecover, &nodeFmtDevNetRecover)
        if err != nil { return err }
        if err := hnm.update910NodeLabel(curNode, newDevRecoverLabel, 
            hnm.getPatchLabel(newNetRecoverSets)); err != nil {
            return err
        }
        lastTimeNetworkRecoverDevices = newNetRecoverSets
    }
    
    return nil
}

3.3 ascendtolerance.go — HotResetTools 深度解析

3.3.1 结构体定义
type HotResetTools struct {
    resetDevNumOnce     int                          // 单次重置设备数(=Ring大小)
    allTaskDevList      map[string][]int32           // 任务名→设备逻辑ID列表
    allTaskDevFaultInfo map[string][]*common.TaskDevInfo // 任务名→设备故障信息列表
    globalDevFaultInfo  map[int32]*common.DevFaultInfo   // 逻辑ID→全局设备故障信息
    taskPod             map[string]v1.Pod            // 任务名→Pod
    faultDev2PodMap     map[int32]v1.Pod             // 故障设备逻辑ID→Pod
    resetTask           map[string]struct{}           // 重置中的任务集合
    resetDev            map[int32]struct{}            // 重置中的设备集合
    queue               workqueue.RateLimitingInterface // K8s工作队列
    podIndexer          cache.Indexer                 // Pod缓存索引器
    cmIndexer           cache.Indexer                 // CM缓存索引器
    jobs                map[string]string             // PodKey→JobName映射
    noResetCmPodKeys    map[string]struct{}           // 无Reset CM的Pod Key集合
}

设计意图

  • resetDevNumOnce 在构造时确定,不同卡类型/用途有不同的Ring大小
  • globalDevFaultInfo 是全局设备故障的单一数据源,每次ListAndWatch周期更新
  • allTaskDevFaultInfoglobalDevFaultInfo派生,加入了Rank信息
  • queue + podIndexer + cmIndexer 构成了K8s Informer事件驱动架构
  • resetTaskresetDev 是两个独立的状态集合,分别跟踪任务级和设备级重置状态
3.3.2 构造函数与Ring大小决策
func NewHotResetManager(devUsage string, deviceNum int, boardId uint32) HotResetManager {
    resetDevNumOnce := getResetDevNumOnce(devUsage, deviceNum, boardId)
    if resetDevNumOnce == 0 {
        return nil  // 不支持的设备类型返回nil
    }
    return &HotResetTools{
        resetDevNumOnce:  resetDevNumOnce,
        resetTask:        map[string]struct{}{},
        resetDev:         map[int32]struct{}{},
        faultDev2PodMap:  map[int32]v1.Pod{},
        jobs:             map[string]string{},
        noResetCmPodKeys: map[string]struct{}{},
    }
}

func getResetDevNumOnce(devUsage string, deviceNum int, boardId uint32) int {
    switch common.ParamOption.RealCardType {
    case api.Ascend910A:
        // 910A: 8卡一个Ring
        return common.Ascend910RingsNum
    case api.Ascend910B:
        if devUsage == common.Infer {
            // 推理卡: A800IA2无HCCS→1卡Ring, 其他→训练Ring
            if boardId == common.A300IA2BoardId || ... {
                return common.Ascend910BRingsNumInfer
            }
            return common.Ascend910BRingsNumTrain
        }
        if devUsage == common.Train {
            // 训练卡: 8卡Ring, 超过8卡→A200T A2 Ring
            resetDevNumOnce = common.Ascend910BRingsNumTrain
            if deviceNum > common.Ascend910BRingsNumTrain {
                return common.A200TA2RingsNum
            }
        }
    case api.Ascend910A3:
        // A3: 全节点关联重置, 设备数=Ring大小
        return deviceNum
    case api.Ascend910A5:
        // A5: 固定Ring大小
        return common.Ascend910A5RingsNum
    }
    return 0  // 不支持的类型
}
3.3.3 核心流程:SyncResetCM — CM/Pod Informer同步

SyncResetCM(ctx, client)

client == nil?

日志错误, return

informers.NewSharedInformerFactory
label: reset=true

Core().V1().ConfigMaps().Informer()

AddEventHandler
(checkConfigMap过滤器)

go cmInformer.Run(ctx.Done())

hrt.queue = client.Queue

hrt.podIndexer = client.PodInformer.GetIndexer()

hrt.cmIndexer = cmInformer.GetIndexer()

WaitForCacheSync
(cmInformer, PodInformer)

go hrt.run()
启动工作队列消费

go func
ctx.Done→queue.ShutDown

完成

CM事件处理

Add/Update

Delete

Event.Type?

handleCMUpdateEvent
更新文件

handleCMDeleteEvent
删除文件

Pod事件处理

Add

Delete

Event.Type?

handlePodAddEvent
创建目录+写CM到文件

handlePodDeleteEvent
删除文件+清理缓存

事件分发 handleEvent

PodResource

CMResource

其他

Event.Resource?

handlePodEvent

handleConfigMapEvent

queue.Forget

工作队列消费循环 run

shutdown

获取obj

run()

processNextWorkItem()

queue.Get()

return false
停止

queue.Done(obj)

是Event?

queue.Forget(obj)

handleEvent(obj)

3.3.4 核心流程:handlePodAddEvent — Pod添加处理

失败

成功

为空

成功

成功

失败

成功

失败

handlePodAddEvent

getPodFromCache
(podIndexer)

重试次数
< MaxPodEventRetryTimes?

queue.AddRateLimited
限速重试

queue.Forget
放弃

GetJobNameOfPod
获取Job名

queue.Forget

hrt.jobs[event.Key] = jobName

writeCmToFileWhilePodAdd

MkdirAll
dataTraceDir

MkdirAll
resetDir

GetCMFromCache
(dataTrace-CM)

writeCMToFile
(写profiling开关)

日志警告

GetCMFromCache
(resetInfo-CM)

writeCMToFile
(写reset信息)

noResetCmPodKeys
标记无Reset CM

queue.Forget
处理完成

设计意图:Pod添加时,为该Pod的Job创建本地目录,并将相关的CM(DataTrace和ResetInfo)写入本地文件。这些文件供设备插件的其他组件(如弹性Agent)读取,实现CM到文件的同步。

3.3.5 核心流程:handleCMUpdateEvent — CM更新处理
func (hrt *HotResetTools) handleCMUpdateEvent(obj interface{}) {
    event := obj.(kubeclient.Event)
    cm, err := hrt.GetCMFromCache(event.Key)
    if err != nil {
        hrt.queue.Forget(obj)
        return
    }
    
    // DataTrace CM: 更新profiling开关文件
    if strings.HasPrefix(cm.Name, common.DataTraceCmPrefix) {
        dir := fmt.Sprintf("%s/%s", common.DataTraceConfigDir, cm.Namespace+"."+cm.Name)
        fileFullName := filepath.Join(dir, common.DataTraceCmProfilingSwitchKey)
        // 只在目录已存在时更新(目录由Pod添加事件创建)
        if _, checkErr := os.Stat(dir); checkErr != nil {
            hrt.queue.Forget(obj)
            return
        }
        hrt.writeCmToFileSystem(cm, common.DataTraceCmProfilingSwitchKey, fileFullName, obj)
        return
    }
    
    // ResetInfo CM: 更新reset信息文件
    if strings.HasPrefix(cm.Name, common.ResetInfoCMNamePrefix) {
        dir := common.GenResetDirName(cm.Namespace, cm.Name)
        if _, checkErr := os.Stat(dir); checkErr != nil {
            hrt.queue.Forget(obj)
            return
        }
        if err = hrt.writeCMToFile(cm); err != nil {
            hrt.queue.AddRateLimited(obj)  // 写入失败→限速重试
            return
        }
    }
    hrt.queue.Forget(obj)
}
3.3.6 核心流程:writeCMToFile — CM数据写入文件

DataTraceCmPrefix

ResetInfoCMNamePrefix

writeCMToFile(cm)

CM名称前缀?

DataTrace处理

目录: DataTraceConfigDir/{ns}.{cm.Name}

读取cm.Data[ProfilingSwitchKey]

WriteToFile
写profiling开关到文件

ResetInfo处理

读取cm.Data[ResetInfoCMDataKey]

WriteToFile
写reset信息到文件
(GenResetFileName)


ResetInfoTypeKey?

WriteToFile
写restartType到文件
(GenResetTypeFileName)

完成

完成

3.3.7 核心流程:GetTaskProcessPolicy — 任务策略决策
func (hrt *HotResetTools) GetTaskProcessPolicy(taskName string) (string, int, error) {
    devFaultInfoList, ok := hrt.allTaskDevFaultInfo[taskName]
    if !ok {
        return "", -1, fmt.Errorf("this task is not in the cache")
    }
    
    var processPolicy string
    var processPolicyLevel int
    // 遍历任务下所有设备的故障信息,取最高级别策略
    for _, devFaultInfo := range devFaultInfoList {
        devPolicyLevel, ok := processPolicyTable[devFaultInfo.Policy]
        if !ok {
            return "", -1, fmt.Errorf("invalid policy of device fault info in task %s", taskName)
        }
        // 取最高级别(最严重)的策略
        if devPolicyLevel > processPolicyLevel {
            processPolicy = devFaultInfo.Policy
            processPolicyLevel = devPolicyLevel
        }
    }
    return processPolicy, processPolicyLevel, nil
}

设计意图:一个任务可能跨多个设备,每个设备有不同的故障级别。任务级策略采用最严重优先原则——以最高级别的设备故障策略作为整个任务的处理策略。

3.3.8 核心流程:UpdateGlobalDevFaultInfoCache — 全局故障缓存更新
func (hrt *HotResetTools) UpdateGlobalDevFaultInfoCache(devDeviceList []*common.NpuDevice, 
    isoDevList []int32) error {
    if len(devDeviceList) == 0 {
        return fmt.Errorf("npu device list is nil")
    }
    // 每次全量重建全局故障缓存
    hrt.globalDevFaultInfo = make(map[int32]*common.DevFaultInfo, len(devDeviceList))
    for _, device := range devDeviceList {
        hrt.globalDevFaultInfo[device.LogicID] = &common.DevFaultInfo{}
        hrt.globalDevFaultInfo[device.LogicID].LogicId = device.LogicID
        hrt.globalDevFaultInfo[device.LogicID].ErrorCode = device.FaultCodes
        
        if common.IntInList(device.LogicID, isoDevList) {
            // 已隔离设备直接设为IsolateError
            hrt.globalDevFaultInfo[device.LogicID].Policy = common.IsolateError
        } else {
            // 根据故障类型获取处理策略
            hrt.globalDevFaultInfo[device.LogicID].Policy = 
                hrt.GetDevProcessPolicy(common.GetFaultType(device.FaultCodes, device.LogicID))
        }
    }
    return nil
}
3.3.9 GetDevProcessPolicy — 设备故障策略映射
func (hrt *HotResetTools) GetDevProcessPolicy(faultType string) string {
    switch faultType {
    case common.NormalNPU, common.NotHandleFault, common.SubHealthFault:
        return common.EmptyError        // Level 0: 无需处理
    case common.RestartRequest:
        return common.RestartRequestError // Level 2: L2请求重启
    case common.RestartBusiness:
        return common.RestartError        // Level 3: L3设备重启
    case common.FreeRestartNPU:
        return common.FreeResetError      // Level 4: L4空闲重置
    case common.RestartNPU:
        return common.ResetError          // Level 5: L5强制重置
    default:
        return common.IsolateError        // Level 6: 隔离
    }
}
3.3.10 核心流程:GenerateTaskDevFaultInfoList — 生成任务设备故障列表
func (hrt *HotResetTools) GenerateTaskDevFaultInfoList(devIdList []int32, 
    rankIndex string) ([]*common.TaskDevInfo, error) {
    // 1. 按LogicID排序,确保Rank分配一致性
    sort.Slice(devIdList, func(i, j int) bool {
        return devIdList[i] < devIdList[j]
    })
    
    // 2. 解析Rank起始索引
    rankStart, err := strconv.Atoi(rankIndex)
    
    devNum := len(devIdList)
    taskDevInfoList := make([]*common.TaskDevInfo, 0, len(devIdList))
    for _, devId := range devIdList {
        var rankId int
        switch rankIndex {
        case common.InferRankIndex:
            // 推理场景:所有设备使用相同的Rank
            rankId = rankStart
        default:
            // 训练场景:Rank = rankStart * devNum + 当前序号
            // 例如:rankStart=0, devNum=8 → 0,1,2,3,4,5,6,7
            //       rankStart=1, devNum=8 → 8,9,10,11,12,13,14,15
            rankId = rankStart*devNum + len(taskDevInfoList)
        }
        
        // 从全局缓存获取设备故障信息
        faultInfo, ok := hrt.globalDevFaultInfo[devId]
        if !ok {
            return nil, fmt.Errorf("device %d is not in global cache", devId)
        }
        
        taskDevInfo := &common.TaskDevInfo{
            RankId:       rankId,
            DevFaultInfo: *faultInfo,  // 值拷贝,避免后续修改影响全局缓存
        }
        taskDevInfoList = append(taskDevInfoList, taskDevInfo)
    }
    return taskDevInfoList, nil
}
3.3.11 核心流程:UpdateFaultDev2PodMap — 故障设备-Pod映射更新
func (hrt *HotResetTools) UpdateFaultDev2PodMap(devList []int32, pod v1.Pod) error {
    for _, device := range devList {
        // 设备有故障→记录映射
        if hrt.globalDevFaultInfo[device].Policy != common.EmptyError &&
            hrt.globalDevFaultInfo[device].Policy != common.IgnoreError {
            hrt.faultDev2PodMap[device] = pod
            continue
        }
        
        // 设备健康但在重置中→不删除映射(重置期间保持关联)
        if _, ok := hrt.resetDev[device]; ok {
            continue
        }
        
        // 设备健康且不在重置中→删除映射
        if _, ok := hrt.faultDev2PodMap[device]; ok {
            delete(hrt.faultDev2PodMap, device)
        }
    }
    return nil
}
3.3.12 核心流程:GetTaskResetInfo — 构建任务重置信息

GetTaskResetInfo
(devFaultInfoList, policy, initPolicy, status)

GetResetDevNumOnce

遍历设备故障信息

策略级别
>= RestartErrorLevel?

跳过(无需重置)

计算Ring索引
= LogicId / resetDevNumOnce

标记Ring为故障Ring

再次遍历设备

设备在
故障Ring中?

跳过

DeepCopyDevInfo
深拷贝设备信息

设置Policy/InitialPolicy/Status

加入rankList

返回TaskResetInfo
{rankList}

设计意图:两阶段处理:

  1. 第一阶段:识别哪些Ring有需要重置的设备(故障级别≥L3)
  2. 第二阶段:将这些Ring上的所有设备(包括非故障设备)都加入重置列表

这是因为Ring级别的重置会影响Ring上所有设备,因此需要将整个Ring的设备状态都更新到CM中。

3.3.13 核心流程:UpdateFreeTask — 清理已结束任务
func (hrt *HotResetTools) UpdateFreeTask(taskListUsedDevice map[string]struct{}, 
    newTaskDevList map[string][]int32) {
    for taskName := range hrt.resetTask {
        // 条件1: 任务不在当前活跃列表中 → 清理
        // 条件2: 任务的设备列表发生变化 → 清理(任务可能被重新调度到其他节点)
        if _, ok := taskListUsedDevice[taskName]; !ok || 
            hrt.isTaskDevListChange(taskName, newTaskDevList) {
            delete(hrt.resetTask, taskName)
        }
    }
}

func (hrt *HotResetTools) isTaskDevListChange(taskName string, 
    newTaskDevList map[string][]int32) bool {
    // 比较旧设备列表和新设备列表是否一致(通过下划线连接后字符串比较)
    return common.Int32Join(hrt.allTaskDevList[taskName], common.UnderLine) !=
        common.Int32Join(newTaskDevList[taskName], common.UnderLine)
}

3.4 ascendcommon_v2.go — A5超节点扩展

3.4.1 文件全貌
// Package device a series of device function
package device

// SetSuperPodType setting the type of super pod
func (tool *AscendTools) SetSuperPodType(superPodType int8) {
    tool.superPodType = superPodType
}

// GetSuperPodType getting the type of super pod
func (tool *AscendTools) GetSuperPodType() int8 {
    return tool.superPodType
}

// SetSuperPodSize setting the type of super pod
func (tool *AscendTools) SetSuperPodSize(superPodSize int32) {
    tool.superPodSize = superPodSize
}

// GetSuperPodSize getting the type of super pod
func (tool *AscendTools) GetSuperPodSize() int32 {
    return tool.superPodSize
}

// SetNodeInternalIPInK8s setting the ip of the node server in k8s
func (tool *AscendTools) SetNodeInternalIPInK8s(nodeIp string) {
    tool.nodeInternalIP = nodeIp
}

// GetNodeInternalIPInK8s getting the ip of the node server in k8s
func (tool *AscendTools) GetNodeInternalIPInK8s() string {
    return tool.nodeInternalIP
}

// SetRackID setting the rank id
func (tool *AscendTools) SetRackID(rackID int32) {
    tool.rackID = rackID
}

// GetRackID getting the rack id
func (tool *AscendTools) GetRackID() int32 {
    return tool.rackID
}

设计意图

  • 这是一个纯getter/setter扩展文件,为AscendTools添加A5超节点相关的属性访问方法
  • 独立成文件的原因:这些属性是A5特有的,与ascendcommon.go中的通用方法逻辑不同,分开便于维护和版本管理
  • 6个方法对应3组属性:superPodType/superPodSize(超Pod规格)、nodeInternalIP(节点IP)、rackID(机架ID)
  • 这些属性在getConfigAnno中被用于构建Pod配置注解的ServerInfo结构
3.4.2 A5属性使用关系

使用场景

ascendcommon_v2.go

SetSuperPodType/GetSuperPodType

SetSuperPodSize/GetSuperPodSize

SetNodeInternalIPInK8s/GetNodeInternalIPInK8s

SetRackID/GetRackID

getConfigAnno
构建ServerInfo

getNodeDeviceInfoCache
A5时写入RackID

外部初始化
设置超Pod属性


四、跨模块协作全景

4.1 ListAndWatch完整周期

DCMI硬件 K8s API HotResetTools HwAscend910Manager AscendTools 主循环 DCMI硬件 K8s API HotResetTools HwAscend910Manager AscendTools 主循环 GetNPUs() GetDeviceList() 设备列表 assemblePhy/VirtualDevices UpdateHealth() GetDeviceAllErrorCode 故障码 isHealthy/isNetworkHealthy setHealthyIfDuoCard/setAICoreHealthyIfVNpu GraceTolerance() UpdateGlobalDevFaultInfoCache setTaskDevInfoCache processAllTask GetTaskProcessPolicy preProcess WriteResetInfoDataIntoCM runProcessTask (goroutine) DoWithVolcanoListAndWatch() getDevStatesDevSet UpdateNodeDeviceInfo WriteDeviceInfoDataIntoCMCache PatchNodeState (Labels) AddPodAnnotation() TryUpdatePodAnnotation HandleDropCardFaultEvents() HandleLostChipFaultEvents() HandleLostNetworkFaultEvents() WriteFaultToEvent (goroutine) CreateEvent

4.2 故障处理完整链路

训练任务 K8s API HotResetTools HwAscend910Manager AscendTools DCMI 故障发生 训练任务 K8s API HotResetTools HwAscend910Manager AscendTools DCMI 故障发生 异步写K8s Event GraceTolerance周期 alt [自愈成功] [自愈失败→升级] alt [L2 (RestartRequest)] [L3 (Restart)] [L5 (Reset)] alt [A3带内重置失败] alt [有任务在用设备] [设备空闲] 硬件故障 故障订阅/轮询获取故障码 writeNewFaultCode isHealthy → Unhealthy allFaultInfo channel CreateEvent(Warning) UpdateGlobalDevFaultInfoCache GetDevProcessPolicy setTaskDevInfoCache GetTaskProcessPolicy restartRequestProcess 检查错误码清除 CM更新为Recovered upgradeRestartRequestProcess 等进程归零→重置设备 CM更新为Recovered/Isolate restartProcess 等进程归零→等自愈→重置 resetProcess 等进程归零→直接重置 hotResetHandler canBeReset → tryResetDevice isRingResetComplete resetDeviceOutBand RescanSoc filterDevStatus (重置中→Healthy) setAllDevUnhealthyOnRing (Ring标记) UpdateNodeDeviceInfo

4.3 Ring级重置拓扑

推理A800IA2 (1卡Ring)

单卡Ring

单卡重置
+阻塞全节点调度

910A5 (超Pod感知)

Ring: A5RingsNum卡

Ring级重置
+超Pod/Rack信息写入CM

910A3 (4设备Ring)

Ring: L0,L1,L2,L3
(2张物理卡×2die)

任一故障→4设备关联重置
带内失败→带外重置
带外失败→第三方扫描

910A/B (8卡Ring)

Ring 0: LogicID 0-7

任一故障→全Ring 8卡重置

Ring 1: LogicID 8-15

任一故障→全Ring 8卡重置

4.4 CM数据模型

本地文件

resetInfo/{ns}.{jobName}/reset_info

resetInfo/{ns}.{jobName}/reset_type

dataTrace/{ns}.{jobName}/profiling_switch

K8s Events

Type: Warning
故障发生事件

Type: Normal
故障恢复事件

Type: Normal
升级原因释放事件

Pod Annotations

AscendReal: 实际使用设备

huawei.com/Ascend910-2kl: klt设备

ascend-910-configuration: 配置JSON

RankIndex: Rank起始索引

Node Labels

Ascend910-Recover: 恢复设备
'0.2.3'

Ascend910-NetworkRecover: 网络恢复设备

server-type: Ascend910-32

server-usage: train/infer

ResetInfo CM

RankList: 重置设备列表
(含LogicID/Policy/Status)

UpdateTime: 更新时间

DeviceInfo CM

Ascend910: 健康设备列表
'Ascend910-0,Ascend910-1'

Ascend910-Recovering: 重置中设备
'Ascend910-2,Ascend910-3'

Ascend910-Unhealthy: 不健康设备

Ascend910-NetworkUnhealthy: 网络不健康

Ascend910-DPUUnhealthy: DPU不健康

Ascend910-FaultList: 故障详情JSON

ManuallySeparateNPU: 手动隔离设备

UpgradeFaultReason: 升级故障原因

SwitchFaultInfo: 交换机故障

DpuInfo: DPU信息


分析总结

这四个文件构成了Ascend Device Plugin的核心设备管理层,实现了以下关键能力:

  1. 设备发现:支持物理设备、虚拟设备(vNPU)、共享设备的统一发现与命名
  2. 健康监测:多维度健康判定(芯片/网络/DPU),支持IPv4/IPv6双栈设备IP获取
  3. 故障容错:六级故障分级(L0-L6),三级升级链(L2→L3→L5→隔离),支持在线/离线两种容错模式
  4. 热重置:Ring级设备重置,A3卡关联重置,带内/带外双重重置路径,第三方设备扫描恢复
  5. CM同步:通过Informer机制实时同步Pod/CM变更,本地文件系统镜像
  6. 调度集成:Volcano调度ListAndWatch,DeviceInfo CM驱动调度决策,Node Label/Annotation管理
  7. A5超Pod:SuperPod/Rack/ServerIndex感知,支持大规模集群拓扑
Logo

鲲鹏昇腾开发者社区是面向全社会开放的“联接全球计算开发者,聚合华为+生态”的社区,内容涵盖鲲鹏、昇腾资源,帮助开发者快速获取所需的知识、经验、软件、工具、算力,支撑开发者易学、好用、成功,成为核心开发者。

更多推荐