mirror of
https://github.com/Awuqing/BackupX.git
synced 2026-09-05 23:47:33 +08:00
功能: v2.2 节点池调度 + Grafana Dashboard + 版本漂移 UI (#49)
节点池动态调度(企业集群核心需求): - model.Node 新增 Labels CSV;Node.HasLabel / LabelSet 辅助方法 - model.BackupTask 新增 NodePoolTag;与 NodeID 互斥(校验层拒绝同时设置) - BackupExecutionService.selectPoolNode:匹配标签的在线节点中选"运行中任务最少" 并列按 ID 升序稳定;空池返回 NODE_POOL_EMPTY 让用户立即感知 - 选中节点仅写 BackupRecord,不回写 task.NodeID —— 每次执行重选实现真轮转均衡 Grafana Dashboard(v2.1 指标的可视化闭环): - deploy/grafana/backupx-dashboard.json:11 个面板覆盖概览/时序/容量/集群 - deploy/grafana/README.md:Prometheus 抓取配置 + 告警建议 - release workflow 打包 grafana/ + nginx.conf 到 tar.gz 前端: - 节点列表:Agent 版本 vs Master 不一致时橙红 Tag + Tooltip 提示升级 - 节点列表新增"标签/节点池"列,支持 CSV 编辑 + 并发/带宽一起改 - 任务表单新增 NodePoolTag 输入框,与节点选择器互斥禁用 测试: - model/node_label_test.go:HasLabel / LabelSet / nil 安全 - service/node_pool_scheduler_test.go:负载最低优先 / 空池错误 / nil repo 降级 - go test ./... + npm run build 全绿
This commit is contained in:
@@ -335,16 +335,29 @@ func (s *BackupExecutionService) startTask(ctx context.Context, id uint, async b
|
||||
nil)
|
||||
}
|
||||
}
|
||||
// 节点池动态选择:task.NodeID=0 且 NodePoolTag 非空时,从匹配的在线节点中挑一台。
|
||||
// 选择策略:正在运行任务数最少者优先;并列时按 ID 升序稳定。
|
||||
// 选中节点仅影响本次运行(task.NodeID 不持久化改动),保证任务在池内轮转。
|
||||
resolvedNodeID := task.NodeID
|
||||
if task.NodeID == 0 && strings.TrimSpace(task.NodePoolTag) != "" {
|
||||
if pooled, perr := s.selectPoolNode(ctx, task.NodePoolTag); perr == nil && pooled != nil {
|
||||
resolvedNodeID = pooled.ID
|
||||
} else if perr != nil {
|
||||
return nil, perr
|
||||
}
|
||||
}
|
||||
startedAt := s.now()
|
||||
// 取第一个存储目标 ID 做兼容
|
||||
primaryTargetID := task.StorageTargetID
|
||||
if tids := collectTargetIDs(task); len(tids) > 0 {
|
||||
primaryTargetID = tids[0]
|
||||
}
|
||||
record := &model.BackupRecord{TaskID: task.ID, StorageTargetID: primaryTargetID, NodeID: task.NodeID, Status: "running", StartedAt: startedAt}
|
||||
record := &model.BackupRecord{TaskID: task.ID, StorageTargetID: primaryTargetID, NodeID: resolvedNodeID, Status: "running", StartedAt: startedAt}
|
||||
if err := s.records.Create(ctx, record); err != nil {
|
||||
return nil, apperror.Internal("BACKUP_RECORD_CREATE_FAILED", "无法创建备份记录", err)
|
||||
}
|
||||
// 用池选出的节点 ID 复写 task 副本,使后续路由/执行沿用
|
||||
task.NodeID = resolvedNodeID
|
||||
task.LastRunAt = &startedAt
|
||||
task.LastStatus = "running"
|
||||
if err := s.tasks.Update(ctx, task); err != nil {
|
||||
@@ -414,6 +427,64 @@ func (s *BackupExecutionService) shouldNotify(ctx context.Context, task *model.B
|
||||
return true
|
||||
}
|
||||
|
||||
// selectPoolNode 从所有 Labels 包含 poolTag 的在线节点中选择"当前运行中任务最少"的一台。
|
||||
// 返回 (nil, error) 表示硬错误(仓储访问失败);(nil, nil) 表示没有匹配节点(退化走本机 Master)。
|
||||
// 本方法不修改任何持久化状态,仅做选择。
|
||||
func (s *BackupExecutionService) selectPoolNode(ctx context.Context, poolTag string) (*model.Node, error) {
|
||||
if s.nodeRepo == nil {
|
||||
// 没接入集群依赖时,降级为让调用方走本机 Master
|
||||
return nil, nil
|
||||
}
|
||||
nodes, err := s.nodeRepo.List(ctx)
|
||||
if err != nil {
|
||||
return nil, apperror.Internal("NODE_LIST_FAILED", "无法枚举节点池", err)
|
||||
}
|
||||
candidates := make([]*model.Node, 0)
|
||||
for i := range nodes {
|
||||
n := &nodes[i]
|
||||
if n.Status != model.NodeStatusOnline {
|
||||
continue
|
||||
}
|
||||
if !n.HasLabel(poolTag) {
|
||||
continue
|
||||
}
|
||||
candidates = append(candidates, n)
|
||||
}
|
||||
if len(candidates) == 0 {
|
||||
return nil, apperror.BadRequest("NODE_POOL_EMPTY",
|
||||
fmt.Sprintf("节点池 %q 下无在线节点,任务无法调度", poolTag), nil)
|
||||
}
|
||||
// 运行中记录数越少越优先。并列按 ID 升序(稳定、可预期)。
|
||||
best := candidates[0]
|
||||
bestLoad := s.countRunningOnNode(ctx, best.ID)
|
||||
for _, n := range candidates[1:] {
|
||||
load := s.countRunningOnNode(ctx, n.ID)
|
||||
if load < bestLoad || (load == bestLoad && n.ID < best.ID) {
|
||||
best = n
|
||||
bestLoad = load
|
||||
}
|
||||
}
|
||||
return best, nil
|
||||
}
|
||||
|
||||
// countRunningOnNode 近似返回节点当前 running 记录数。失败按 0 处理(不影响功能,仅退化调度精度)。
|
||||
func (s *BackupExecutionService) countRunningOnNode(ctx context.Context, nodeID uint) int {
|
||||
if s.records == nil {
|
||||
return 0
|
||||
}
|
||||
items, err := s.records.List(ctx, repository.BackupRecordListOptions{Status: model.BackupRecordStatusRunning})
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
count := 0
|
||||
for i := range items {
|
||||
if items[i].NodeID == nodeID {
|
||||
count++
|
||||
}
|
||||
}
|
||||
return count
|
||||
}
|
||||
|
||||
// effectiveBandwidth 返回当前上下文应用的带宽限速字符串。
|
||||
// 优先级:Node.BandwidthLimit(非空) > 全局 s.bandwidthLimit。
|
||||
func (s *BackupExecutionService) effectiveBandwidth(ctx context.Context, nodeID uint) string {
|
||||
|
||||
Reference in New Issue
Block a user