学习目标
类比:运维自动化就像写字楼的「中央工单机器人」。过去保洁、检修要挨个房间敲门(人工登每台集群),现在把任务写成「剧本」(Job 模板),到点或审批通过后,机器人经统一门禁(K8s API 网关)进对应房间干活;每干一次都留痕(审计),敏感活儿还要主管签字(审批)。
学完本章你应该能够:
- 区分定时任务(CronJob)、自动化编排(Playbook / Job 模板)、跨集群批量运维三类能力的适用场景。
- 在统一观测底座上,把自动化任务的「执行耗时 / 成功率 / 失败原因」作为指标纳入 Prometheus 联邦监控(标签严格为
tenant/app/env/cluster)。 - 配置跨集群批量运维任务,任务经统一 K8s API 网关下发到目标物理隔离集群执行,且执行受 RBAC 与 Namespace 租户隔离约束。
- 为敏感 / 高危自动化任务配置多级审批与执行审计留痕,满足金融 / 政企不可绕过要求。
- 在测试环境与生产环境之间区分任务并发、审批门槛、失败熔断与回收策略。
- 动手写出可编译的 Go 调度器:用
gin+client-go实现 Cron 调度、Argo Workflows 提交、经 API 网关跨集群执行、审批钩子、任务指标上报(复用模块 05 联邦)。
前置知识:
- 已掌握《模块 02 多 K8s 集群纳管》(统一 API 网关、物理隔离、Namespace 租户隔离)。
- 已掌握《模块 05 统一观测平台》(联邦指标、标签体系、Alertmanager 分级告警)。
- 熟悉 K8s 原生
CronJob/Job,了解 Argo Workflows 或等价编排引擎。 - 了解多级审批流(模块 11)与审计日志(模块 09)的基本概念。
- Go 基础:
ginWeb 框架、client-go操作 K8s、gRPC/HTTP 客户端。
本章你会动手做的事:
- 写一个每天凌晨清理测试集群闲置资源的
CronJob,并配置并发控制与超时。 - 把一个「跨三集群批量扩容」操作建成 Job 模板,提交生产审批后执行。
- 编译运行 Go 调度器,提交一个跨集群任务,观察指标
task_executions_total、task_duration_seconds进入联邦监控。
一、模块概述与企业商用价值
1.1 这个模块解决什么痛点
在多租户、多物理隔离集群的底座上,运维有两个真实困境:
- 重复劳动:扩缩容、配置下发、巡检要在 N 套集群上重复 N 次,人工易错且不可追溯。
- 失控风险:谁在什么时候对哪台集群做了什么,若没有审批与留痕,金融合规直接不达标;一个高危脚本误执行可能跨租户影响他人。
运维自动化与定时任务模块用「编排 + 审批 + 审计 + 联邦可观测」把零散操作变成受控、可追、可观测的标准动作,是平台「降本 + 合规」的交汇点。
1.2 平台不可替代性
- 自助运维必须「经审批触发」:租户点一下按钮,背后是多级审批 + 网关下发 + 执行留痕,杜绝越权。
- 跨集群批量动作若直连各集群 apiserver,等于绕过统一门禁;本模块强制所有执行经 API 网关,保留 RBAC 与审计边界。
- 任务执行数据进联邦监控,失败率 / 耗时异常能自动告警,把「自动化翻车」变成「可预警」。
1.3 模块在 PaaS 全局中的定位
下图说明:自动化模块位于「人 / 策略」与「多集群资源」之间,所有执行经 API 网关,所有结果进联邦监控,所有敏感动作受审批 / 审计约束。
graph TD A[运维自动化与定时任务模块] --> B[统一 K8s API 网关] A --> C[联邦监控体系] A --> D[平台标准管控] B --> B1[物理隔离 集群A] B --> B2[物理隔离 集群B] B --> B3[租户 Namespace 隔离] C --> C1[任务指标标签] C --> C2[Thanos 长期存储] D --> D1[RBAC 细粒度] D --> D2[审计 90 天+] D --> D3[多级审批]
这张图在讲什么:自动化不是「绕过管控的快进键」,而是把管控(网关 / RBAC / 审批 / 审计)编织进每一次自动执行的通道里。
二、细分功能详解(商用生产级)
下面按基础能力 / 高级企业增值能力分级。★禁止删减项:审批、审计、权限隔离——它们是自动化模块合规的底线,缺任一即为不合格。
2.1 功能清单对照
| 能力分级 | 功能项 | 说明 | 商用要点 |
|---|---|---|---|
| 基础能力 | ① 定时任务 | Cron 表达式、时区、并发控制 | 周期运维 |
| 基础能力 | ② 自动化编排 | Playbook / Job 模板 | 可复用动作 |
| 基础能力 | ③ 自助运维 | 经审批触发 | ★禁止删减:审批 |
| 高级企业增值 | ④ 跨集群批量运维 | 扩缩容 / 配置下发 / 巡检 | ★禁止删减:权限隔离 |
| 高级企业增值 | ⑤ 任务审计与执行留痕 | 谁/何时/在哪集群/结果 | ★禁止删减:审计 |
| 高级企业增值 | ⑥ 失败重试与熔断 | 退避重试 / 熔断阈值 | 防雪崩 |
| 高级企业增值 | ⑦ 资源回收与闲置治理 | 识别并回收闲置 | 降本 |
2.2 基础能力逐条拆解
① 定时任务(CronJob)
基于 K8s 原生 CronJob,支持标准 5 段 Cron、指定时区(spec.timeZone)、并发策略(Forbid / Replace)、起始最后期限(startingDeadlineSeconds)与超时(activeDeadlineSeconds)。例如「每日 02:00 清理测试集群 Evicted Pod」。
② 自动化编排(Playbook / Job 模板)
把「登录 → 检查 → 变更 → 校验」封装成模板,支持参数化(目标集群、Namespace、副本数)。底层优先用 Argo Workflows 或等价 DAG 引擎表达依赖;简单单步用 Job 模板即可。
③ 自助运维(★禁止删减:审批)
租户在控制台点「执行」,请求进入《模块 11 审批流》:低风险(如测试环境重启 Pod)可自动放行;高风险(生产扩缩容、配置下发)必须多级审批,审批通过才由网关下发执行。审批不可绕过是金融底线。
2.3 高级企业增值能力
④ 跨集群批量运维(★禁止删减:权限隔离)
一次操作对多套物理隔离集群生效(如三集群统一下发 ConfigMap、统一巡检)。任务经 API 网关逐集群下发,每个目标集群按 RBAC 校验 Namespace 租户权限——A 租户的任务绝不允许触达 B 租户 Namespace。权限隔离是跨集群批量动作的安全闸门。
⑤ 任务审计与执行留痕(★禁止删减:审计)
每次任务记录:提交人、审批人、目标集群、Namespace、参数、开始/结束时间、退出码、输出摘要,写入审计日志留存 ≥ 90 天。结合模块 09 全局审计视图可回溯任意历史执行。
⑥ 失败重试与熔断
支持指数退避重试;当单任务失败率超阈值或目标集群不可达,自动熔断并告警,避免重试风暴拖垮网关。
⑦ 资源回收与闲置治理
定期扫描低利用率工作负载(基于联邦指标),生成回收建议工单,执行同样走审批 + 审计。
下图把上述能力落到「触发—编排—执行—治理」四层。
graph LR
subgraph 触发层
T1[定时 CronJob]
T2[自助运维申请]
end
subgraph 编排层
P[Playbook / Job 模板]
W[Argo Workflows DAG]
end
subgraph 执行层
GW[统一 K8s API 网关]
NS[租户 Namespace 隔离]
end
subgraph 治理层
AU[审计留痕 90天+]
AL[联邦指标告警]
end
T1 --> P
T2 --> W
P --> GW
W --> GW
GW --> NS
NS --> AU
NS --> AL这张图在讲什么:从触发到治理,任务每一步都被网关收口、被 Namespace 隔离、被审计与监控覆盖。
2.4 项目结构与文件清单
本章所有代码统一收口为一个可独立构建的 Go 服务 paas-automation,并配合一套 K8s IaC 清单。目录树如下(文件名与后文逐个展示一一对应):
paas-automation/
├── go.mod # Go 依赖(gin / client-go / argo / cron / prometheus)
├── cmd/
│ └── scheduler/
│ └── main.go # gin 服务入口:任务提交/查询 API
├── internal/
│ ├── auth/
│ │ └── tenant.go # 复用模块05:租户/审批身份中间件
│ ├── scheduler/
│ │ └── cron.go # Cron 定时调度器
│ ├── workflow/
│ │ └── argo.go # Argo Workflows 提交(client-go)
│ ├── gateway/
│ │ └── executor.go # 经 API 网关跨集群执行(client-go)
│ ├── approval/
│ │ └── hook.go # 多级审批钩子
│ └── metrics/
│ └── reporter.go # 任务指标上报(tenant/app/env/cluster)
└── deploy/ # K8s IaC(见第三、四章)
├── 00-namespace-rbac.yaml # automation ns + 跨集群执行 SA/ClusterRole
├── 01-cronjob-test.yaml # 测试环境 CronJob(清闲置)
├── 02-cronjob-prod.yaml # 生产环境 CronJob(差异化)
├── 03-job-template-scaleout.yaml # 跨集群批量扩容 Job 模板
├── 04-argo-workflow-batch.yaml # Argo Workflows 批量编排
├── 05-workflow-crd.yaml # Workflow CRD 引用说明
└── 06-task-metrics.yaml # 任务指标 ServiceMonitor + recording rules
后文按「架构联动(第三章)→ 端到端 SOP(第四章)→ 管控与故障(第五、六章)」逐文件展开。Go 上报的指标名(
task_executions_total/task_duration_seconds)与 4.7 recording rules、模块 05 的task:success_rate:rate7d完全一致。
三、底层架构联动设计
3.1 与多 K8s 集群 / API 网关的交互
所有自动化任务的执行体(Job / Workflow Pod)不在平台控制面直接调目标集群 apiserver,而是经统一 K8s API 网关反向隧道下发。网关侧按 RBAC 解析「该任务对目标集群 / 命名空间是否有写权限」,校验通过才转发;kubeconfig 仅在网关侧经 KMS 信封解密,执行端不落盘明文。这保证物理隔离不被自动化绕开。
3.2 ★ 跨集群执行联动图(必画)
下图展示:任务在平台侧生成 → 审批通过 → 网关逐集群下发 → 目标集群按租户 Namespace 隔离执行 → 执行指标带 tenant/app/env/cluster 进联邦。
graph LR
subgraph Plat[平台控制面]
J[Job / Workflow 模板]
AP[审批流引擎]
SCH[Go 调度器]
end
subgraph Gate[统一 K8s API 网关]
RB[RBAC 校验]
KMS[KMS 解密 kubeconfig]
end
subgraph C1[生产集群A]
NS1[租户 Namespace]
end
subgraph C2[测试集群B]
NS2[租户 Namespace]
end
FM[(联邦 Prometheus)]
SCH --> J
J --> AP
AP --> RB
RB --> KMS
KMS --> NS1
KMS --> NS2
NS1 -- 执行指标带标签 --> FM
NS2 -- 执行指标带标签 --> FM这张图在讲什么:审批是闸口、RBAC+KMS 是门锁、Namespace 是房间、联邦指标是留痕——四者共同保证「自动但不失控」。
三个关键设计点:
- 网关收口:跨集群执行必须经网关,杜绝执行端直连,保留统一审计与权限边界。
- 标签贯穿:任务执行指标(耗时、成功/失败、目标集群)注入
tenant/app/env/cluster,在联邦侧即可按租户看自动化健康度(Go 见 4.9)。 - 与告警/审批/审计联动:执行失败率超阈值触发 Alertmanager 分级告警(复用模块 05);每次执行事件进审计;敏感任务审批事件是放行的唯一凭证。
3.3 与 RBAC / 审批 / 审计的联动
- RBAC 细粒度:任务模板绑定「可操作集群 / Namespace / 动作」三元组;越权模板在提交阶段即被拒。
- 多级审批:生产环境高危模板(扩缩容、配置下发、删除类)必须多级审批,审批不通过网关拒绝下发。
- 审计 90 天+:提交、审批、下发、执行结果全链路写入审计,结合模块 05 联邦监控可同时看「做了什么」与「运行得怎样」。
3.4 K8s IaC(一):automation 命名空间 + 跨集群执行 RBAC
执行端 ServiceAccount 仅授予目标租户命名空间的「最小权限」,绝不绑 cluster-admin。
# file: deploy/00-namespace-rbac.yaml
apiVersion: v1
kind: Namespace
metadata:
name: automation
labels:
tenant: platform
env: production
cluster: central
---
apiVersion: v1
kind: ServiceAccount
metadata:
name: automation-runner
namespace: automation
---
# 仅允许操作 tn-acme 命名空间内的 deployments(限副本数变更/只读),禁止 secrets/删除
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: automation-runner-tenant-acme
namespace: tn-acme
rules:
- apiGroups: ["apps"]
resources: ["deployments"]
verbs: ["get", "list", "patch", "update"]
- apiGroups: [""]
resources: ["pods", "pods/log"]
verbs: ["get", "list"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: automation-runner-tenant-acme
namespace: tn-acme
roleRef:
apiGroup: rbac.authorization.k8s.io
kind: Role
name: automation-runner-tenant-acme
subjects:
- kind: ServiceAccount
name: automation-runner
namespace: automation
---
# 平台控制面调用网关的权限(仅向网关发起,不直接触达集群)
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: gateway-invoker
rules:
- apiGroups: [""]
resources: ["services"]
resourceNames: ["k8s-api-gateway"]
verbs: ["get", "create"]
3.5 K8s IaC(二):CronJob 模板(测试 / 生产差异化)
# file: deploy/01-cronjob-test.yaml
# 【测试环境】低风险,免审批,直接执行;并发宽松
apiVersion: batch/v1
kind: CronJob
metadata:
name: cleanup-evicted-test
namespace: automation
labels:
tenant: platform
env: testing
cluster: testing-cluster
spec:
schedule: "0 2 * * *"
timeZone: "Asia/Shanghai"
concurrencyPolicy: Forbid
startingDeadlineSeconds: 300
jobTemplate:
spec:
activeDeadlineSeconds: 600
backoffLimit: 2
template:
spec:
serviceAccountName: automation-runner
containers:
- name: cleanup
image: registry.internal/kubectl:1.29
command: ["sh", "-c", "kubectl delete pod -n tn-acme --field-selector status.phase=Failed"]
restartPolicy: Never
# file: deploy/02-cronjob-prod.yaml
# 【生产环境】差异化:限流、严格时窗、失败即告警;提交需审批
apiVersion: batch/v1
kind: CronJob
metadata:
name: cleanup-evicted-prod
namespace: automation
labels:
tenant: platform
env: production
cluster: prod-cluster-a
spec:
schedule: "0 3 * * *"
timeZone: "Asia/Shanghai"
concurrencyPolicy: Forbid
startingDeadlineSeconds: 200
successfulJobsHistoryLimit: 3
failedJobsHistoryLimit: 3
jobTemplate:
spec:
activeDeadlineSeconds: 900
backoffLimit: 1 # 生产失败不盲目重试,转人工
template:
metadata:
labels:
tenant: platform
env: production
cluster: prod-cluster-a
spec:
serviceAccountName: automation-runner
containers:
- name: cleanup
image: registry.internal/kubectl:1.29
command: ["sh", "-c", "kubectl delete pod -n tn-acme --field-selector status.phase=Failed"]
restartPolicy: Never
⚠️
concurrencyPolicy: Forbid防止上一次没跑完又触发;serviceAccountName必须只有目标租户 Namespace 的删改权限,绝不能绑 cluster-admin。生产环境backoffLimit设小,避免重试风暴(见 3.3 熔断)。
3.6 K8s IaC(三):Job 模板(跨集群批量扩容)
# file: deploy/03-job-template-scaleout.yaml
# 跨集群批量扩容:由 Go 调度器按目标集群分别渲染并提交(见 4.8)
apiVersion: batch/v1
kind: Job
metadata:
name: scale-out-tn-acme # 实际由调度器加集群后缀:scale-out-tn-acme-prod-a
namespace: automation
labels:
tenant: acme
app: web
env: production
cluster: prod-cluster-a
task-template: scale-out
spec:
backoffLimit: 1
activeDeadlineSeconds: 300
template:
spec:
serviceAccountName: automation-runner
restartPolicy: Never
containers:
- name: scale
image: registry.internal/kubectl:1.29
command:
- sh
- -c
- |
kubectl -n tn-acme patch deployment web --type merge \
-p "{\"spec\":{\"replicas\":${REPLICAS}}}"
env:
- name: REPLICAS
value: "6"
3.7 K8s IaC(四):Argo Workflows 批量编排 + CRD
# file: deploy/04-argo-workflow-batch.yaml
# Argo Workflows:跨集群巡检 DAG(先 prod-a,再 prod-b,最后汇总共识)
apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
name: batch-inspect-acme
namespace: automation
labels:
tenant: acme
env: production
task-template: batch-inspect
spec:
entrypoint: main
templates:
- name: main
dag:
tasks:
- name: inspect-prod-a
template: inspect
arguments: {parameters: [{name: cluster, value: prod-cluster-a}]}
- name: inspect-prod-b
template: inspect
dependencies: [inspect-prod-a]
arguments: {parameters: [{name: cluster, value: prod-cluster-b}]}
- name: inspect
inputs:
parameters:
- name: cluster
container:
image: registry.internal/kubectl:1.29
command: ["sh", "-c", "kubectl --context {{inputs.parameters.cluster}} get pods -n tn-acme"]
# file: deploy/05-workflow-crd.yaml
# Workflow CRD 由 argo 安装提供,此处为平台纳管说明(不重复定义,仅引用)
# 安装:kubectl apply -n argo -f https://github.com/argoproj/argo-workflows/releases/download/v3.5.4/install.yaml
# 平台侧仅创建 Workflow 实例(见 04),并绑定 RBAC 到 automation-runner。
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
metadata:
name: workflows.argoproj.io
spec:
group: argoproj.io
names:
kind: Workflow
plural: workflows
singular: workflow
scope: Namespaced
versions:
- name: v1alpha1
served: true
storage: true
3.8 K8s IaC(五):任务指标 ServiceMonitor + recording rules
把任务执行指标纳入模块 05 联邦。指标名与标签与 Go 上报(4.9)严格一致。
# file: deploy/06-task-metrics.yaml
apiVersion: v1
kind: ConfigMap
metadata:
name: task-recording-rules
namespace: monitoring
data:
task_recording.yml: |
groups:
- name: task_recording
interval: 30s
rules:
# 近 7 天成功率(供 Grafana / 告警复用)
- record: task:success_rate:rate7d
expr: |
sum by (tenant, env, cluster) (rate(task_executions_total{result="success"}[7d]))
/
sum by (tenant, env, cluster) (rate(task_executions_total[7d]))
# 平均耗时
- record: task:duration_avg:5m
expr: sum by (tenant, env, cluster) (rate(task_duration_seconds_sum[5m])) / sum by (tenant, env, cluster) (rate(task_duration_seconds_count[5m]))
- name: task_alerts
rules:
# 跨集群执行失败率超阈值 → 触发分级告警(复用模块05 Alertmanager)
- alert: TaskFailureRateHigh
expr: |
sum by (tenant, env, cluster) (rate(task_executions_total{result="failed"}[1h]))
/
sum by (tenant, env, cluster) (rate(task_executions_total[1h])) > 0.2
for: 10m
labels:
severity: warning
annotations:
summary: "租户 {{ $labels.tenant }} 任务失败率超 20%"
---
# 任务执行 Pod 暴露 /metrics(pushgateway 或 sidecar 暴露),由 Agent 抓取
apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
name: task-executor
namespace: monitoring
labels:
tenant: platform
cluster: central
spec:
selector:
matchLabels:
app: paas-automation-scheduler
endpoints:
- port: metrics
interval: 30s
四、端到端标准操作流程
按角色拆分,并显式区分测试环境与生产环境差异。
4.1 角色与职责
- 开发(租户):提交自助运维申请,查看本租户任务执行状态。
- 集群运维:编写 / 维护 Job 模板与 CronJob,处理跨集群批量动作。
- 平台管理员:配置审批策略、RBAC 边界、审计规则与熔断阈值。
4.2 标准操作流程(SOP)
步骤 1|创建任务模板(集群运维)
- 步骤 1.1 编写 Job / Workflow 模板,声明目标集群、Namespace、参数(见 3.6 / 3.7)。
- 步骤 1.2 绑定 RBAC 三元组(集群 / Namespace / 动作),限定可操作范围(防越权,见 3.4)。
步骤 2|提交与审批(开发 / 平台管理员)
- 步骤 2.1 **【测试环境】**低风险任务可直接执行,免审批。
- 步骤 2.2 **【生产环境】**必须走多级审批;审批通过才生成网关下发凭证(Go 见 4.7
approval.Hook)。
步骤 3|网关下发与执行(平台自动)
- 步骤 3.1 网关校验 RBAC + KMS 解密,逐集群下发(Go 见 4.8
gateway.Executor)。 - 步骤 3.2 目标集群在租户 Namespace 内执行,指标带标签进联邦(Go 见 4.9
metrics.Report)。
步骤 4|观测与留痕(全员)
- 步骤 4.1 Grafana 看任务成功率 / 耗时(recording rules 见 3.8);失败走 Alertmanager 分级。
- 步骤 4.2 执行详情进审计,留存 ≥ 90 天。
4.3 定时任务提交示例(平台 CLI 经网关)
# 提交一个跨集群批量扩容任务,经网关下发(生产必走审批)
paas-cli task submit \
--template scale-out \
--clusters prod-a,prod-b,test-b \
--namespace tn-acme \
--params replicas=6 \
--approve-required true # 生产必走审批
4.4 Go 调度器:gin 服务入口
平台 Go 服务 paas-automation 用 gin 暴露任务提交 / 查询 API,串联审批钩子、跨集群执行器、指标上报。下面对应 2.4 目录树各文件。
// file: cmd/scheduler/main.go
package main
import (
"context"
"fmt"
"log"
"net/http"
"os"
"strings"
"time"
"paas-automation/internal/approval"
"paas-automation/internal/auth"
"paas-automation/internal/gateway"
"paas-automation/internal/metrics"
"paas-automation/internal/scheduler"
"github.com/gin-gonic/gin"
)
func main() {
gwAddr := os.Getenv("API_GATEWAY_ADDR")
if gwAddr == "" {
gwAddr = "http://k8s-api-gateway.platform.svc:8080"
}
approvalAddr := os.Getenv("APPROVAL_ADDR")
if approvalAddr == "" {
approvalAddr = "http://approval-svc.platform.svc:8080"
}
exec := gateway.NewExecutor(gwAddr)
approver := approval.NewHook(approvalAddr)
cron := scheduler.NewCronScheduler()
store := scheduler.NewMemoryStore()
api := &apiHandler{
exec: exec,
approver: approver,
cron: cron,
store: store,
}
r := gin.Default()
r.Use(auth.AuthTenant()) // 复用模块05 风格:解析 X-Tenant / X-Role
v1 := r.Group("/api/v1")
{
// 提交一次性跨集群任务
v1.POST("/tasks", api.submitTask)
// 注册定时任务
v1.POST("/schedules", api.addSchedule)
// 查询任务状态
v1.GET("/tasks/:id", api.getTask)
// 指标自上报端点(供联邦抓取,标签含 tenant/app/env/cluster)
v1.GET("/metrics", gin.WrapH(metrics.PromHandler()))
}
cron.Start()
if err := r.Run(getenv("LISTEN_ADDR", ":8090")); err != nil {
log.Fatalf("scheduler exit: %v", err)
}
}
type apiHandler struct {
exec *gateway.Executor
approver *approval.Hook
cron *scheduler.CronScheduler
store *scheduler.MemoryStore
}
// getenv 读取环境变量,缺失时返回默认值。
func getenv(k, def string) string {
if v := os.Getenv(k); v != "" {
return v
}
return def
}
// runTask 执行一次跨集群任务:审批 → 网关下发 → 指标上报。
func (h *apiHandler) runTask(ctx context.Context, t scheduler.TaskRequest) error {
approved, err := h.approver.RequireApproval(ctx, approval.Request{
Tenant: t.Tenant, Cluster: t.Cluster, Namespace: t.Namespace, Action: t.Action, Env: t.Env,
})
if err != nil {
return err
}
if !approved {
return fmt.Errorf("approval required or rejected")
}
start := time.Now()
_, err = h.exec.ExecuteOnCluster(ctx, gateway.ExecSpec{
Cluster: t.Cluster, Namespace: t.Namespace, Tenant: t.Tenant, Env: t.Env, App: t.App, Manifest: t.Manifest,
})
elapsed := time.Since(start).Seconds()
outcome := "success"
if err != nil {
outcome = "failed"
}
metrics.Report(t.Tenant, t.App, t.Env, t.Cluster, t.Template, outcome, elapsed)
return err
}
// submitTask 提交一次性任务:先经审批钩子,再经网关逐集群下发。
func (h *apiHandler) submitTask(c *gin.Context) {
var req scheduler.TaskRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if err := h.runTask(c.Request.Context(), req); err != nil {
if strings.Contains(err.Error(), "approval") {
c.JSON(http.StatusForbidden, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{"status": "dispatched", "task_id": req.ID})
}
// addSchedule 注册定时任务,到点在 runTask 中执行(同样走审批/网关/指标)。
func (h *apiHandler) addSchedule(c *gin.Context) {
var req scheduler.TaskRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if _, err := h.cron.Add(req, func(ctx context.Context, t scheduler.TaskRequest) error {
return h.runTask(ctx, t)
}); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
h.store.Add(req)
c.JSON(http.StatusOK, gin.H{"status": "scheduled"})
}
// getTask 查询任务详情(示例存储)。
func (h *apiHandler) getTask(c *gin.Context) {
t, ok := h.store.Get(c.Param("id"))
if !ok {
c.JSON(http.StatusNotFound, gin.H{"error": "task not found"})
return
}
c.JSON(http.StatusOK, t)
}
上例
authTenant()、time等辅助在 4.10、4.11 给出,保证整段可编译。
4.5 Go:Cron 定时调度器
// file: internal/scheduler/cron.go
package scheduler
import (
"context"
"sync"
"time"
"github.com/robfig/cron/v3"
)
// TaskRequest 为提交任务的请求体;字段与 3.8 指标标签一致。
type TaskRequest struct {
ID string `json:"id"`
Tenant string `json:"tenant"`
App string `json:"app"`
Env string `json:"env"`
Cluster string `json:"cluster"`
Namespace string `json:"namespace"`
Template string `json:"template"`
Action string `json:"action"`
Manifest string `json:"manifest"`
CronExpr string `json:"cron_expr"`
}
// MemoryStore 为示例存储(生产替换为 etcd/数据库)。
type MemoryStore struct {
mu sync.RWMutex
tasks map[string]TaskRequest
}
func NewMemoryStore() *MemoryStore {
return &MemoryStore{tasks: map[string]TaskRequest{}}
}
// Add 写入一条任务(生产替换为 etcd/数据库)。
func (m *MemoryStore) Add(t TaskRequest) {
m.mu.Lock()
defer m.mu.Unlock()
m.tasks[t.ID] = t
}
// Get 按 ID 读取任务。
func (m *MemoryStore) Get(id string) (TaskRequest, bool) {
m.mu.RLock()
defer m.mu.RUnlock()
t, ok := m.tasks[id]
return t, ok
}
// CronScheduler 封装 robfig/cron 调度。
type CronScheduler struct {
c *cron.Cron
}
func NewCronScheduler() *CronScheduler {
return &CronScheduler{c: cron.New(cron.WithLocation(time.Local))}
}
// Add 注册一个定时任务;回调里触发跨集群执行(通过 ExecuteFunc 注入执行器)。
func (s *CronScheduler) Add(t TaskRequest, run func(ctx context.Context, t TaskRequest) error) (cron.EntryID, error) {
return s.c.AddFunc(t.CronExpr, func() {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
defer cancel()
_ = run(ctx, t) // 实际执行由 main 注入,含审批/网关/指标
})
}
func (s *CronScheduler) Start() {
s.c.Start()
}
func (s *CronScheduler) Stop() {
s.c.Stop()
}
4.6 Go:Argo Workflows 提交(client-go)
// file: internal/workflow/argo.go
package workflow
import (
"context"
wfv1alpha1 "github.com/argoproj/argo-workflows/v3/pkg/apis/workflow/v1alpha1"
wfclientset "github.com/argoproj/argo-workflows/v3/pkg/client/clientset/versioned"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
// Submit 通过 Argo client-go 提交一个 Workflow 实例到 automation 命名空间。
// 注意:真正实现中该 client 的 rest.Config 由 API 网关经 KMS 解密后注入(见 gateway.Executor)。
func Submit(ctx context.Context, client wfclientset.Interface, ns string, wf *wfv1alpha1.Workflow) (*wfv1alpha1.Workflow, error) {
return client.ArgoprojV1alpha1().Workflows(ns).Create(ctx, wf, metav1.CreateOptions{})
}
// NewBatchInspect 构造一个跨集群巡检 Workflow(与 deploy/04 对应)。
func NewBatchInspect(name, tenant string) *wfv1alpha1.Workflow {
return &wfv1alpha1.Workflow{
ObjectMeta: metav1.ObjectMeta{
Name: name,
Namespace: "automation",
Labels: map[string]string{
"tenant": tenant,
"env": "production",
"task-template": "batch-inspect",
},
},
Spec: wfv1alpha1.WorkflowSpec{
Entrypoint: "main",
Templates: []wfv1alpha1.Template{
{
Name: "main",
DAG: &wfv1alpha1.DAGTemplate{
Tasks: []wfv1alpha1.DAGTask{
{Name: "inspect-prod-a", Template: "inspect",
Arguments: wfv1alpha1.Arguments{Parameters: []wfv1alpha1.Parameter{{Name: "cluster", Value: strPtr("prod-cluster-a")}}}},
{Name: "inspect-prod-b", Template: "inspect", Dependencies: []string{"inspect-prod-a"},
Arguments: wfv1alpha1.Arguments{Parameters: []wfv1alpha1.Parameter{{Name: "cluster", Value: strPtr("prod-cluster-b")}}}},
},
},
},
{
Name: "inspect",
Inputs: wfv1alpha1.Inputs{Parameters: []wfv1alpha1.Parameter{{Name: "cluster"}}},
Container: &wfv1alpha1.Container{
Image: "registry.internal/kubectl:1.29",
Command: []string{"sh", "-c", "kubectl --context {{inputs.parameters.cluster}} get pods -n tn-acme"},
},
},
},
},
}
}
func strPtr(s string) *string { return &s }
4.7 Go:审批钩子
// file: internal/approval/hook.go
package approval
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
)
// Hook 调用平台审批服务,实现多级审批闸口。
type Hook struct {
BaseURL string
HTTP *http.Client
}
// NewHook 构造审批钩子客户端。
func NewHook(baseURL string) *Hook {
return &Hook{BaseURL: baseURL, HTTP: &http.Client{Timeout: 15 * time.Second}}
}
// Request 与任务提交请求对齐。
type Request struct {
Tenant string `json:"tenant"`
Cluster string `json:"cluster"`
Namespace string `json:"namespace"`
Action string `json:"action"`
Env string `json:"env"`
}
// RequireApproval 向审批服务询问是否放行。
// 测试环境低风险可由审批服务直接自动通过;生产高风险须人工多级审批。
func (h *Hook) RequireApproval(ctx context.Context, req Request) (bool, error) {
body, err := json.Marshal(req)
if err != nil {
return false, err
}
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost,
h.BaseURL+"/v1/approve", bytes.NewReader(body))
if err != nil {
return false, err
}
httpReq.Header.Set("Content-Type", "application/json")
resp, err := h.HTTP.Do(httpReq)
if err != nil {
return false, err
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
b, _ := io.ReadAll(resp.Body)
return false, fmt.Errorf("approval http %d: %s", resp.StatusCode, string(b))
}
var out struct {
Approved bool `json:"approved"`
}
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return false, err
}
return out.Approved, nil
}
4.8 Go:经 API 网关跨集群执行器(client-go)
执行器不直接连目标集群,而是把 manifest 经统一 K8s API 网关下发;网关侧做 RBAC + KMS 解密。
// file: internal/gateway/executor.go
package gateway
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
)
// ExecSpec 描述一次跨集群执行;标签 tenant/app/env/cluster 用于指标与审计。
type ExecSpec struct {
Tenant string `json:"tenant"`
App string `json:"app"`
Env string `json:"env"`
Cluster string `json:"cluster"`
Namespace string `json:"namespace"`
Manifest string `json:"manifest"`
}
// Result 为网关返回的执行结果。
type Result struct {
TaskID string `json:"task_id"`
Status string `json:"status"`
}
// Executor 通过统一 API 网关下发执行(不直接连接目标集群 apiserver)。
type Executor struct {
GatewayBaseURL string
HTTP *http.Client
}
// NewExecutor 构造网关执行器。
func NewExecutor(gatewayBaseURL string) *Executor {
return &Executor{
GatewayBaseURL: gatewayBaseURL,
HTTP: &http.Client{Timeout: 60 * time.Second},
}
}
// ExecuteOnCluster 把 manifest 经网关下发到指定集群的目标租户 Namespace。
// 网关侧会:① RBAC 校验 Namespace 租户权限;② KMS 解密该集群 kubeconfig;
// ③ 在 tn-<tenant> 内 apply manifest;④ 返回 task_id。
// 权限隔离在此闸口强制生效:A 租户任务无法触达 B 租户 Namespace。
func (e *Executor) ExecuteOnCluster(ctx context.Context, spec ExecSpec) (*Result, error) {
body, err := json.Marshal(spec)
if err != nil {
return nil, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
e.GatewayBaseURL+"/v1/cluster/"+spec.Cluster+"/apply", bytes.NewReader(body))
if err != nil {
return nil, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("X-Tenant", spec.Tenant) // 网关据此做租户级 RBAC
req.Header.Set("X-Env", spec.Env)
resp, err := e.HTTP.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode >= 300 {
b, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("gateway apply http %d: %s", resp.StatusCode, string(b))
}
var out Result
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
return nil, err
}
return &out, nil
}
4.9 Go:任务指标上报(tenant/app/env/cluster)
使用 prometheus/client_golang,指标名与 3.8 recording rules、模块 05 task:success_rate:rate7d 完全对齐。
// file: internal/metrics/reporter.go
package metrics
import (
"net/http"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
var (
// 执行次数计数器:标签严格为 tenant/app/env/cluster/template/result
executions = promauto.NewCounterVec(prometheus.CounterOpts{
Name: "task_executions_total",
Help: "Total automation task executions, labeled by tenant/app/env/cluster.",
}, []string{"tenant", "app", "env", "cluster", "template", "result"})
// 执行耗时直方图(秒)
duration = promauto.NewHistogramVec(prometheus.HistogramOpts{
Name: "task_duration_seconds",
Help: "Automation task duration in seconds.",
Buckets: []float64{1, 5, 15, 30, 60, 120, 300},
}, []string{"tenant", "app", "env", "cluster", "template"})
)
// Report 上报一次任务执行结果(供联邦抓取,标签与 recording rules 一致)。
func Report(tenant, app, env, cluster, template, result string, durSeconds float64) {
executions.WithLabelValues(tenant, app, env, cluster, template, result).Inc()
duration.WithLabelValues(tenant, app, env, cluster, template).Observe(durSeconds)
}
// PromHandler 暴露 /metrics 端点,供 Prometheus Agent 抓取并联邦上报。
func PromHandler() http.Handler {
return promhttp.Handler()
}
4.10 Go:租户/审批身份中间件(复用模块05 风格)
// file: internal/auth/tenant.go
package auth
import (
"github.com/gin-gonic/gin"
)
// AuthTenant 解析 X-Tenant / X-Role,供调度器鉴权与审计使用。
func AuthTenant() gin.HandlerFunc {
return func(c *gin.Context) {
tenant := c.GetHeader("X-Tenant")
if tenant == "" {
tenant = "unknown"
}
role := c.GetHeader("X-Role") // platform-admin / tenant
c.Set("tenant", tenant)
c.Set("role", role)
c.Next()
}
}
4.11 Go 依赖与构建
// file: go.mod
module paas-automation
go 1.21
require (
github.com/gin-gonic/gin v1.9.1
github.com/argoproj/argo-workflows/v3 v3.5.4
github.com/robfig/cron/v3 v3.0.1
github.com/prometheus/client_golang v1.18.0
k8s.io/client-go v0.29.0
k8s.io/apimachinery v0.29.0
)
cd paas-automation
go mod tidy
go build ./...
go run ./cmd/scheduler
# 提交一个生产跨集群扩容任务(审批通过后才下发)
curl -XPOST http://localhost:8090/api/v1/tasks -H 'Content-Type: application/json' -d '{
"tenant":"acme","app":"web","env":"production","cluster":"prod-cluster-a",
"namespace":"tn-acme","template":"scale-out","action":"patch",
"manifest":"apiVersion: batch/v1\nkind: Job\n..."
}'
4.12 端到端时序图
sequenceDiagram participant Dev as 开发(租户) participant SCH as Go 调度器 participant AP as 审批服务 participant GW as 统一 K8s API 网关 participant CA as 目标集群 apiserver participant FM as 联邦 Prometheus Dev->>SCH: 提交扩容申请(跨3集群) SCH->>AP: RequireApproval AP-->>SCH: 多级审批通过 SCH->>GW: ExecuteOnCluster(经网关) GW->>CA: RBAC校验+KMS解密后下发 CA->>FM: 执行指标带 tenant 标签 FM-->>Dev: Grafana 展示任务成功率
这张图在讲什么:从申请到看到结果,审批是闸口、网关是通道、联邦是留痕,三者缺一不可。
五、生产环境管控与安全约束
5.1 资源配额与权限隔离
- 自动化执行 Pod 运行在独立
automationNamespace,与业务租户 Namespace 隔离;其 ServiceAccount 经 RBAC 仅授予目标租户范围的最小权限(见 3.4)。 - 跨集群批量任务强制逐集群 RBAC 校验,单任务不允许跨租户 Namespace 操作(网关侧
X-Tenant强制)。 - 高危动作(delete / scale 到 0 / 配置覆盖)标记为敏感,必须二次认证 + 审批。
5.2 敏感操作审批与审计
- 生产环境所有「写 / 删 / 扩缩容 / 配置下发」类任务走多级审批(模块 11),审批事件不可绕过、不可篡改(Go 见 4.7
approval.Hook)。 - 每次执行记录提交人、审批人、目标集群、Namespace、参数、结果、耗时,写入审计日志留存 ≥ 90 天。
- 任务模板的增改同样走审批 + 审计,防止「恶意模板」入库。
5.3 故障熔断策略
- 目标集群经网关不可达:任务标记失败并熔断该集群后续下发,触发 Alertmanager 分级告警(recording rule
TaskFailureRateHigh)。 - 单任务失败率超阈值:自动暂停同模板后续实例,进入人工确认。
- 重试采用指数退避(
backoffLimit+ 间隔),避免雪崩(见 3.5 生产backoffLimit: 1)。
5.4 测试环境 vs 生产环境 差异化管控对照表
| 管控项 | 测试环境 | 生产环境 |
|---|---|---|
| 审批门槛 | 低风险免审批 | 写/删/扩缩容必多级审批 |
| 并发执行 | 可高并发 | 限流 + 分批,防影响业务 |
| 权限边界 | 限本测试集群 | 严格 RBAC,禁跨租户 |
| 失败处理 | 记录即可 | 熔断 + 告警升级 |
| 审计留存 | 30 天 | ≥ 90 天(合规) |
| 模板变更 | 运维直接改 | 审批 + 审计 + 二次认证 |
| 资源回收 | 可直接清 | 生成工单,审批后执行 |
| CronJob backoff | 可重试 2 次 | 限 1 次,转人工 |
六、常见生产故障与解决方案
6.1 高频故障清单
故障 1|跨集群批量任务部分集群失败
- 现象:3 集群中 1 个失败,其余成功。
- 排查:多为该集群网关 Endpoint 摘除或 RBAC 缺权限。
- 优化:先在单集群灰度验证模板;网关侧自动摘除异常集群并告警。
故障 2|CronJob 漏执行
- 现象:某天定时任务没跑。
- 排查:
startingDeadlineSeconds过小导致错过窗口,或控制面调度拥塞。 - 优化:放宽 deadline、错峰调度、加「未执行」告警。
故障 3|自动化误删业务 Pod
- 现象:租户反馈服务中断。
- 排查:模板 ServiceAccount 权限过大,越权触达业务 Namespace。
- 优化:收紧 RBAC 到租户级(见 3.4);删除类强制审批 + 二次认证;执行前 dry-run 校验。
故障 4|重试风暴拖垮网关
- 现象:网关 CPU 飙升,正常运维卡顿。
- 排查:
backoffLimit过高且无退避间隔。 - 优化:指数退避 + 熔断阈值;单模板并发上限。
6.2 运维 runbook(真实命令)
注册 / 校验跨集群执行权限
# 1. 校验 automation-runner 对 tn-acme 的真实权限(绝不应含 secrets/delete 全局)
kubectl -n tn-acme auth can-i patch deployments --as=system:serviceaccount:automation:automation-runner
kubectl -n tn-acme auth can-i delete secrets --as=system:serviceaccount:automation:automation-runner
# 期望:第一条 yes,第二条 no
# 2. 经网关下发一个 Job,观察 task_id(网关侧做 RBAC+KMS)
curl -XPOST http://k8s-api-gateway.platform.svc:8080/v1/cluster/prod-cluster-a/apply \
-H "X-Tenant: acme" -H "X-Env: production" -H "Content-Type: application/json" \
-d @job-scaleout.json
触发 / 校验任务指标进入联邦
# 调度器 /metrics 暴露指标,确认命名空间 tn-acme 的 CronJob 是否上报
kubectl -n automation port-forward svc/paas-automation 8090:8090 &
curl -s 'http://localhost:8090/metrics' | grep '^task_executions_total'
# 在中心 Thanos Query 跨集群查该租户成功率(复用模块05)
curl -s 'http://thanos-query.monitoring.svc:9090/api/v1/query?query=task:success_rate:rate7d{tenant="acme"}' | jq .
# 注册 ServiceMonitor 让 Agent 抓取调度器指标(见 deploy/06)
kubectl -n monitoring apply -f deploy/06-task-metrics.yaml
审批流验证(生产必走)
# 直接提交生产高危任务应被拒(未审批)
curl -XPOST http://localhost:8090/api/v1/tasks -H 'Content-Type: application/json' -d '{
"tenant":"acme","env":"production","cluster":"prod-cluster-a",
"namespace":"tn-acme","template":"scale-out","action":"patch","manifest":"..."}'
# 期望返回 403 approval required or rejected
# 走审批通过后重试应成功(审批服务返回 approved=true)
查审计 / 失败率告警
# 查审计中该租户的任务执行记录(模块09 审计视图)
paas-cli audit query --tenant acme --last 7d
# 失败率超阈值时,Alertmanager 应产生 TaskFailureRateHigh(复用模块05 分级路由)
kubectl -n monitoring port-forward svc/alertmanager 9093:9093 &
curl 'http://localhost:9093/api/v2/alerts' | jq '.[] | select(.labels.alertname=="TaskFailureRateHigh")'
6.3 故障排查流程图
flowchart TD
S[任务执行异常告警] --> C{失败范围}
C -->|单集群| T1[检查网关 Endpoint 与 RBAC]
C -->|全部| T2[检查模板与审批状态]
C -->|定时漏跑| T3[检查 deadline 与调度]
T1 --> R1[重发或摘除异常集群]
T2 --> R2[复核权限与审批]
T3 --> R3[放宽窗口并加未执行告警]
R1 --> E[恢复]
R2 --> E
R3 --> E这张图在讲什么:按失败范围二分定位——单集群查网关/RBAC,全失败查模板/审批,漏跑查调度,快速收敛。
自测题与动手练习
5 道自测题
- 为什么跨集群批量运维必须经由统一 K8s API 网关下发,而不能让执行端直连各集群 apiserver?
- ★ 本模块三项禁止删减能力(审批 / 审计 / 权限隔离)各自防止了哪类风险?
CronJob的concurrencyPolicy、startingDeadlineSeconds、activeDeadlineSeconds分别控制什么?- 自动化任务的执行指标如何进入 Prometheus 联邦?需要带哪些标签?(提示:见 4.9 与 3.8)
- 测试环境与生产环境在「审批门槛」和「审计留存」上有哪些差异化管控?Go 调度器如何实现审批闸口?
3 个动手练习
- 写一个每天 03:00 清理
tn-demo命名空间 Failed Pod 的CronJob,配置Forbid并发与 600s 超时,并约束 ServiceAccount 仅限该命名空间(参考 3.4 / 3.5)。 - 把一个「跨 prod-a / prod-b 批量下发 ConfigMap」操作建成 Job 模板,列出它需要绑定的 RBAC 三元组与必须走的审批级别(参考 3.4 / 4.7)。
- 在 Grafana 用联邦指标写一条 PromQL:按
tenant聚合近 7 天自动化任务成功率sum by (tenant) (rate(task_executions_total{result="success"}[7d])) / sum by (tenant) (rate(task_executions_total[7d]))(与 recording ruletask:success_rate:rate7d等价)。
本章小结
- 运维自动化与定时任务模块把重复运维变成「受控、可追、可观测」的标准动作,是平台降本与合规的交汇点。
- ★ 三件不可删减的事:审批(敏感动作不可绕过,Go 见
approval.Hook)、审计(执行全链路留痕 ≥ 90 天)、权限隔离(跨集群执行受 RBAC + Namespace 租户边界约束,Go 见gateway.Executor的X-Tenant闸口)。 - 所有跨集群执行经统一 K8s API 网关收口,配合 KMS 解密 kubeconfig 与 RBAC 校验,物理隔离不被自动化绕开。
- 任务执行指标带
tenant/app/env/cluster进联邦监控(Gometrics.Report+task_executions_total/task_duration_seconds,recording ruletask:success_rate:rate7d在模块 05 已复用),失败率 / 耗时异常联动 Alertmanager 分级告警,把「自动化翻车」变为「可预警」。 - 测试 / 生产在审批门槛、并发、熔断、审计留存、模板变更上必须差异化(对照表 5.4);下一篇《模块 07 安全管控》将进一步用镜像扫描与零信任加固这些自动化通道。