Ansible 管配置、Helm 管部署、ArgoCD 管同步、CronJob 管定时、GitLab CI 管流水线——但谁来管它们?用一套自建巡检系统把所有工具串起来,7 分钟跑完全量检查,异常自动告警。
早上到公司,打开 6 个浏览器标签页:Prometheus 看监控、Grafana 看仪表盘、ArgoCD 看同步状态、GitLab 看 Pipeline 结果、K9s 看 Pod 状态、Ansible Tower 看批量任务。花了 20 分钟确认系统正常,然后重复了 365 天。
这不是运维,这是保安巡逻。前 6 篇已经教你把各个工具搭起来,这篇解决的是:怎么把它们串成一套自动化巡检系统——7 分钟跑完全量检查,异常自动推送到手机。
为什么要自建巡检系统
你可能会说:Prometheus + AlertManager 不是已经能告警了吗?
对,但它只覆盖了指标层面。下面这些场景它管不了:
| 场景 |
Prometheus 能发现吗 |
巡检系统能发现吗 |
| 某台机器的 SSH 端口不通 |
❌(机器挂了,没指标了) |
✅ |
| ArgoCD Application 同步失败 |
❌ |
✅ |
| GitLab Pipeline 失败没人管 |
❌ |
✅ |
| 证书 7 天后到期 |
⚠️(需要额外 Exporter) |
✅ |
| Helm 部署的 Chart 版本过旧 |
❌ |
✅ |
| 磁盘 inode 耗尽 |
⚠️(需要 Node Exporter) |
✅ |
| Ansible 批量任务执行失败 |
❌ |
✅ |
| 某个 CronJob 连续 3 次失败 |
❌ |
✅ |
所以巡检系统的定位不是替代监控,而是覆盖监控的盲区:配置漂移、CI/CD 状态、定时任务健康、基础设施连通性。监控管不着的角落,交给巡检系统兜底。
架构设计
┌─────────────────────────────────────────────────────┐
│ 巡检系统架构 │
├─────────────────────────────────────────────────────┤
│ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Ansible │ │ K8s API │ │ GitLab │ │
│ │ Playbook │ │ 查询 │ │ API │ │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌────────────────────────────────────────┐ │
│ │ 巡检引擎 (Python) │ │
│ │ ┌──────┐ ┌──────┐ ┌──────┐ ┌──────┐ │ │
│ │ │检查器1│ │检查器2│ │检查器3│ │检查器N│ │ │
│ │ └──────┘ └──────┘ └──────┘ └──────┘ │ │
│ └──────────────────┬─────────────────────┘ │
│ │ │
│ ┌──────┴──────┐ │
│ │ 结果聚合 │ │
│ └──────┬──────┘ │
│ │ │
│ ┌──────────┼──────────┐ │
│ ▼ ▼ ▼ │
│ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │Web 仪表盘│ │企业微信 │ │邮件告警 │ │
│ └────────┘ └────────┘ └────────┘ │
│ │
└─────────────────────────────────────────────────────┘
核心思路很直接:每个检查项是一个独立的 Python 函数,返回统一的结果格式。巡检引擎负责调度、聚合、告警,检查项之间完全解耦。
巡检引擎设计
项目结构
inspection-system/
├── engine/
│ ├── __init__.py
│ ├── runner.py # 巡检引擎
│ ├── reporter.py # 结果聚合
│ ├── notifier.py # 告警通知
│ └── checks/
│ ├── __init__.py
│ ├── ssh_check.py # SSH 连通性
│ ├── disk_check.py # 磁盘空间
│ ├── cert_check.py # 证书到期
│ ├── k8s_check.py # K8s 集群健康
│ ├── argocd_check.py # ArgoCD 同步状态
│ ├── gitlab_check.py # GitLab Pipeline 状态
│ ├── helm_check.py # Helm Release 版本
│ └── cronjob_check.py # CronJob 执行状态
├── config/
│ ├── targets.yaml # 巡检目标配置
│ └── alert_rules.yaml # 告警规则
├── templates/
│ └── report.html # Web 仪表盘模板
├── app.py # Web 入口
├── run_inspection.py # 命令行入口
└── Dockerfile
检查结果统一格式
每个检查器都返回相同的数据结构,聚合起来就很方便:
# engine/checks/base.py
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
class Severity(Enum):
OK = "ok"
WARNING = "warning"
CRITICAL = "critical"
UNKNOWN = "unknown"
@dataclass
class CheckResult:
name: str # 检查项名称
target: str # 检查目标
severity: Severity # 严重级别
message: str # 结果描述
details: dict = field(default_factory=dict) # 详细信息
timestamp: datetime = field(default_factory=datetime.now)
@property
def passed(self):
return self.severity == Severity.OK
def to_dict(self):
return {
"name": self.name,
"target": self.target,
"severity": self.severity.value,
"message": self.message,
"details": self.details,
"timestamp": self.timestamp.isoformat(),
}
巡检引擎
# engine/runner.py
import importlib
import concurrent.futures
from datetime import datetime
from typing import List
from .checks.base import CheckResult, Severity
class InspectionRunner:
def __init__(self, config: dict):
self.config = config
self.results: List[CheckResult] = []
def _load_checks(self) -> list:
"""根据配置动态加载检查器"""
checks = []
for check_config in self.config.get("checks", []):
module_name = f"engine.checks.{check_config['module']}"
module = importlib.import_module(module_name)
check_class = getattr(module, check_config["class"])
checks.append(check_class(check_config.get("params", {})))
return checks
def run_single(self, check) -> CheckResult:
"""执行单个检查器,捕获异常避免一个崩了全崩"""
try:
return check.run()
except Exception as e:
return CheckResult(
name=check.__class__.__name__,
target="engine",
severity=Severity.UNKNOWN,
message=f"检查器执行异常: {e}",
)
def run_all(self) -> List[CheckResult]:
"""并发执行所有检查器"""
checks = self._load_checks()
self.results = []
with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:
futures = {executor.submit(self.run_single, check): check for check in checks}
for future in concurrent.futures.as_completed(futures):
result = future.result()
self.results.append(result)
# 实时打印结果
icon = {"ok": "✅", "warning": "⚠️", "critical": "🔴", "unknown": "❓"}
print(f"{icon.get(result.severity.value, '❓')} [{result.severity.value.upper():>8}] "
f"{result.name:>20} | {result.target:>30} | {result.message}")
return self.results
def summary(self) -> dict:
"""生成汇总报告"""
total = len(self.results)
by_severity = {}
for sev in Severity:
by_severity[sev.value] = sum(1 for r in self.results if r.severity == sev)
return {
"total": total,
"by_severity": by_severity,
"healthy": by_severity.get("ok", 0),
"issues": total - by_severity.get("ok", 0),
"timestamp": datetime.now().isoformat(),
"results": [r.to_dict() for r in self.results],
}
核心检查器
检查器 1:SSH 连通性
# engine/checks/ssh_check.py
import paramiko
from .base import CheckResult, Severity
class SSHCheck:
"""检查目标机器的 SSH 连通性"""
def __init__(self, params: dict):
self.hosts = params.get("hosts", [])
self.port = params.get("port", 22)
self.timeout = params.get("timeout", 10)
self.key_file = params.get("key_file", "~/.ssh/id_rsa")
def run(self) -> CheckResult:
failed_hosts = []
for host in self.hosts:
try:
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(host, port=self.port, timeout=self.timeout,
key_filename=self.key_file)
client.close()
except Exception as e:
failed_hosts.append({"host": host, "error": str(e)})
if not failed_hosts:
return CheckResult(
name="SSH连通性",
target=f"{len(self.hosts)}台主机",
severity=Severity.OK,
message=f"全部 {len(self.hosts)} 台主机 SSH 可达",
)
elif len(failed_hosts) < len(self.hosts):
return CheckResult(
name="SSH连通性",
target=f"{len(failed_hosts)}/{len(self.hosts)}台不可达",
severity=Severity.WARNING,
message=f"{len(failed_hosts)} 台主机 SSH 不可达",
details={"failed_hosts": failed_hosts},
)
else:
return CheckResult(
name="SSH连通性",
target=f"{len(self.hosts)}台主机",
severity=Severity.CRITICAL,
message="所有主机 SSH 不可达",
details={"failed_hosts": failed_hosts},
)
检查器 2:磁盘空间
# engine/checks/disk_check.py
import paramiko
from .base import CheckResult, Severity
class DiskCheck:
"""通过 SSH 检查磁盘空间使用率"""
def __init__(self, params: dict):
self.hosts = params.get("hosts", [])
self.warning_threshold = params.get("warning", 80)
self.critical_threshold = params.get("critical", 90)
self.key_file = params.get("key_file", "~/.ssh/id_rsa")
def _check_host(self, host: str) -> dict:
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
client.connect(host, key_filename=self.key_file, timeout=10)
# df -h 输出解析,排除 tmpfs/overlay 等伪文件系统
stdin, stdout, stderr = client.exec_command(
"df -h --type=ext4 --type=xfs --type=btrfs 2>/dev/null || df -h | grep -vE 'tmpfs|overlay|udev'"
)
output = stdout.read().decode()
client.close()
disks = []
for line in output.strip().split("\n")[1:]:
parts = line.split()
if len(parts) >= 6:
mount = parts[5]
usage_str = parts[4].replace("%", "")
try:
usage = int(usage_str)
disks.append({"mount": mount, "usage": usage, "total": parts[1], "used": parts[2]})
except ValueError:
continue
return {"host": host, "disks": disks}
def run(self) -> CheckResult:
results = []
warnings = []
criticals = []
for host in self.hosts:
try:
result = self._check_host(host)
results.append(result)
for disk in result["disks"]:
if disk["usage"] >= self.critical_threshold:
criticals.append(f"{host}:{disk['mount']} ({disk['usage']}%)")
elif disk["usage"] >= self.warning_threshold:
warnings.append(f"{host}:{disk['mount']} ({disk['usage']}%)")
except Exception as e:
criticals.append(f"{host}: 检查失败 ({e})")
if criticals:
return CheckResult(
name="磁盘空间",
target=f"{len(self.hosts)}台主机",
severity=Severity.CRITICAL,
message=f"{len(criticals)} 个磁盘超过 {self.critical_threshold}%",
details={"criticals": criticals, "warnings": warnings},
)
elif warnings:
return CheckResult(
name="磁盘空间",
target=f"{len(self.hosts)}台主机",
severity=Severity.WARNING,
message=f"{len(warnings)} 个磁盘超过 {self.warning_threshold}%",
details={"warnings": warnings},
)
else:
return CheckResult(
name="磁盘空间",
target=f"{len(self.hosts)}台主机",
severity=Severity.OK,
message="所有磁盘使用率正常",
)
检查器 3:K8s 集群健康
# engine/checks/k8s_check.py
from kubernetes import client, config
from .base import CheckResult, Severity
class K8sHealthCheck:
"""检查 K8s 集群健康状态"""
def __init__(self, params: dict):
self.kubeconfig = params.get("kubeconfig", "~/.kube/config")
self.namespace = params.get("namespace", "")
def run(self) -> CheckResult:
config.load_kube_config(self.kubeconfig)
v1 = client.CoreV1Api()
apps_v1 = client.AppsV1Api()
batch_v1 = client.BatchV1Api()
issues = []
# 1. 检查节点状态
nodes = v1.list_node()
for node in nodes.items:
for condition in node.status.conditions:
if condition.type == "Ready" and condition.status != "True":
issues.append(f"节点 {node.metadata.name} NotReady")
# 2. 检查异常 Pod
all_ns = "" if self.namespace else True
pods = v1.list_pod_for_all_namespaces() if not self.namespace else v1.list_namespaced_pod(self.namespace)
abnormal_pods = []
for pod in pods.items:
if pod.status.phase not in ("Running", "Succeeded"):
abnormal_pods.append(f"{pod.metadata.namespace}/{pod.metadata.name} ({pod.status.phase})")
if abnormal_pods:
issues.append(f"异常 Pod {len(abnormal_pods)} 个: {abnormal_pods[:5]}")
# 3. 检查失败的 CronJob
cronjobs = batch_v1.list_cron_job_for_all_namespaces() if not self.namespace else batch_v1.list_namespaced_cron_job(self.namespace)
failed_cronjobs = []
for cj in cronjobs.items:
if cj.status.last_schedule_time:
# 检查最近一次 Job 是否失败
jobs = batch_v1.list_namespaced_job(cj.metadata.namespace)
for job in jobs.items:
if (job.metadata.owner_references and
job.metadata.owner_references[0].name == cj.metadata.name):
if job.status.failed and job.status.failed > 0:
failed_cronjobs.append(f"{cj.metadata.namespace}/{cj.metadata.name}")
break
if failed_cronjobs:
issues.append(f"失败的 CronJob: {failed_cronjobs}")
node_count = len(nodes.items)
ready_count = sum(1 for n in nodes.items
if any(c.type == "Ready" and c.status == "True" for c in n.status.conditions))
if issues:
severity = Severity.CRITICAL if any("NotReady" in i for i in issues) else Severity.WARNING
return CheckResult(
name="K8s集群健康",
target=f"{ready_count}/{node_count} 节点就绪",
severity=severity,
message=f"发现 {len(issues)} 个问题",
details={"issues": issues},
)
else:
return CheckResult(
name="K8s集群健康",
target=f"{node_count} 节点",
severity=Severity.OK,
message=f"全部 {node_count} 个节点就绪,无异常 Pod",
)
检查器 4:ArgoCD 同步状态
# engine/checks/argocd_check.py
import requests
from .base import CheckResult, Severity
class ArgoCDCheck:
"""检查 ArgoCD Application 同步状态"""
def __init__(self, params: dict):
self.url = params.get("url", "https://argocd.homelab.local")
self.token = params.get("token", "")
def run(self) -> CheckResult:
headers = {"Authorization": f"Bearer {self.token}"}
resp = requests.get(f"{self.url}/api/v1/applications",
headers=headers, verify=False, timeout=10)
resp.raise_for_status()
apps = resp.json().get("items", [])
out_of_sync = []
degraded = []
for app in apps:
name = app["metadata"]["name"]
sync_status = app.get("status", {}).get("sync", {}).get("status", "Unknown")
health_status = app.get("status", {}).get("health", {}).get("status", "Unknown")
if sync_status != "Synced":
out_of_sync.append(f"{name} (sync={sync_status})")
if health_status not in ("Healthy", "Progressing"):
degraded.append(f"{name} (health={health_status})")
if degraded:
return CheckResult(
name="ArgoCD同步",
target=f"{len(apps)}个应用",
severity=Severity.CRITICAL,
message=f"{len(degraded)} 个应用健康状态异常",
details={"degraded": degraded, "out_of_sync": out_of_sync},
)
elif out_of_sync:
return CheckResult(
name="ArgoCD同步",
target=f"{len(apps)}个应用",
severity=Severity.WARNING,
message=f"{len(out_of_sync)} 个应用未同步",
details={"out_of_sync": out_of_sync},
)
else:
return CheckResult(
name="ArgoCD同步",
target=f"{len(apps)}个应用",
severity=Severity.OK,
message=f"全部 {len(apps)} 个应用已同步且健康",
)
检查器 5:GitLab Pipeline 状态
# engine/checks/gitlab_check.py
import requests
from datetime import datetime, timedelta
from .base import CheckResult, Severity
class GitLabPipelineCheck:
"""检查 GitLab 最近 Pipeline 执行状态"""
def __init__(self, params: dict):
self.url = params.get("url", "https://gitlab.homelab.local")
self.token = params.get("token", "")
self.project_ids = params.get("project_ids", [])
self.lookback_hours = params.get("lookback_hours", 24)
def run(self) -> CheckResult:
headers = {"PRIVATE-TOKEN": self.token}
since = (datetime.utcnow() - timedelta(hours=self.lookback_hours)).isoformat()
failed_pipelines = []
total_pipelines = 0
for project_id in self.project_ids:
resp = requests.get(
f"{self.url}/api/v4/projects/{project_id}/pipelines",
headers=headers,
params={"updated_after": since, "per_page": 50},
timeout=10,
)
resp.raise_for_status()
pipelines = resp.json()
for p in pipelines:
total_pipelines += 1
if p["status"] == "failed":
# 获取项目名称
proj_resp = requests.get(
f"{self.url}/api/v4/projects/{project_id}",
headers=headers, timeout=10,
)
proj_name = proj_resp.json().get("path_with_namespace", str(project_id))
failed_pipelines.append({
"project": proj_name,
"branch": p["ref"],
"url": f"{self.url}/{proj_name}/-/pipelines/{p['id']}",
})
if failed_pipelines:
return CheckResult(
name="GitLab Pipeline",
target=f"{len(self.project_ids)}个项目",
severity=Severity.WARNING,
message=f"最近 {self.lookback_hours}h 内 {len(failed_pipelines)}/{total_pipelines} 个 Pipeline 失败",
details={"failed_pipelines": failed_pipelines},
)
else:
return CheckResult(
name="GitLab Pipeline",
target=f"{len(self.project_ids)}个项目",
severity=Severity.OK,
message=f"最近 {self.lookback_hours}h 内 {total_pipelines} 个 Pipeline 全部成功",
)
检查器 6:证书到期检查
# engine/checks/cert_check.py
import ssl
import socket
from datetime import datetime
from .base import CheckResult, Severity
class CertificateCheck:
"""检查 HTTPS 证书到期时间"""
def __init__(self, params: dict):
self.domains = params.get("domains", [])
self.warning_days = params.get("warning_days", 30)
self.critical_days = params.get("critical_days", 7)
def _check_cert(self, domain: str, port: int = 443) -> dict:
context = ssl.create_default_context()
context.check_hostname = False
context.verify_mode = ssl.CERT_NONE
with socket.create_connection((domain, port), timeout=10) as sock:
with context.wrap_socket(sock, server_hostname=domain) as ssock:
cert = ssock.getpeercert()
# 证书格式: 'Jul 15 23:59:59 2026 GMT'
expire_str = cert.get("notAfter", "")
expire_date = datetime.strptime(expire_str, "%b %d %H:%M:%S %Y %Z")
days_left = (expire_date - datetime.utcnow()).days
return {"domain": domain, "expire_date": expire_str, "days_left": days_left}
def run(self) -> CheckResult:
results = []
warnings = []
criticals = []
for domain in self.domains:
try:
result = self._check_cert(domain)
results.append(result)
if result["days_left"] <= self.critical_days:
criticals.append(f"{domain} ({result['days_left']}天)")
elif result["days_left"] <= self.warning_days:
warnings.append(f"{domain} ({result['days_left']}天)")
except Exception as e:
criticals.append(f"{domain}: 检查失败 ({e})")
if criticals:
return CheckResult(
name="证书到期",
target=f"{len(self.domains)}个域名",
severity=Severity.CRITICAL,
message=f"{len(criticals)} 个证书 {self.critical_days} 天内到期",
details={"criticals": criticals, "warnings": warnings},
)
elif warnings:
return CheckResult(
name="证书到期",
target=f"{len(self.domains)}个域名",
severity=Severity.WARNING,
message=f"{len(warnings)} 个证书 {self.warning_days} 天内到期",
details={"warnings": warnings},
)
else:
return CheckResult(
name="证书到期",
target=f"{len(self.domains)}个域名",
severity=Severity.OK,
message="所有证书有效期充足",
details={"all": results},
)
告警通知
通知模块
# engine/notifier.py
import requests
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from typing import List
from .checks.base import CheckResult, Severity
class Notifier:
def __init__(self, config: dict):
self.wecom_webhook = config.get("wecom_webhook", "")
self.smtp_config = config.get("smtp", {})
self.alert_recipients = config.get("recipients", [])
def send(self, results: List[CheckResult]):
"""只发送非 OK 的结果"""
issues = [r for r in results if r.severity != Severity.OK]
if not issues:
return
# 企业微信通知
if self.wecom_webhook:
self._send_wecom(issues)
# 邮件通知
if self.smtp_config and self.alert_recipients:
self._send_email(issues)
def _send_wecom(self, issues: List[CheckResult]):
"""企业微信群机器人 Webhook"""
lines = ["## 运维巡检告警\n"]
for r in issues:
icon = {"warning": "⚠️", "critical": "🔴", "unknown": "❓"}
lines.append(f"{icon.get(r.severity.value, '❓')} **{r.name}** | {r.target}")
lines.append(f" {r.message}\n")
payload = {
"msgtype": "markdown",
"markdown": {"content": "\n".join(lines)}
}
requests.post(self.wecom_webhook, json=payload, timeout=10)
def _send_email(self, issues: List[CheckResult]):
"""邮件通知"""
msg = MIMEMultipart("alternative")
msg["Subject"] = f"[运维巡检] 发现 {len(issues)} 个异常"
msg["From"] = self.smtp_config.get("from", "ops@homelab.local")
msg["To"] = ", ".join(self.alert_recipients)
html_parts = ["<html><body><h2>运维巡检报告</h2><table border='1' cellpadding='8'>"]
html_parts.append("<tr><th>级别</th><th>检查项</th><th>目标</th><th>描述</th></tr>")
for r in issues:
color = {"warning": "#FFA500", "critical": "#FF0000", "unknown": "#808080"}
html_parts.append(
f"<tr><td style='color:{color.get(r.severity.value, '#000')}'>{r.severity.value.upper()}</td>"
f"<td>{r.name}</td><td>{r.target}</td><td>{r.message}</td></tr>"
)
html_parts.append("</table></body></html>")
msg.attach(MIMEText("".join(html_parts), "html", "utf-8"))
with smtplib.SMTP(self.smtp_config["host"], self.smtp_config["port"]) as server:
server.starttls()
server.login(self.smtp_config["user"], self.smtp_config["password"])
server.send_message(msg)
配置文件
巡检目标配置
# config/targets.yaml
checks:
# SSH 连通性
- module: ssh_check
class: SSHCheck
params:
hosts:
- 192.168.1.101 # k3s-master
- 192.168.1.102 # k3s-node1
- 192.168.1.103 # k3s-node2
- 192.168.1.104 # dell-3430
- 192.168.1.105 # synology
port: 22
timeout: 10
# 磁盘空间
- module: disk_check
class: DiskCheck
params:
hosts:
- 192.168.1.101
- 192.168.1.102
- 192.168.1.103
- 192.168.1.104
- 192.168.1.105
warning: 80
critical: 90
# K8s 集群健康
- module: k8s_check
class: K8sHealthCheck
params:
kubeconfig: ~/.kube/config
# ArgoCD 同步状态
- module: argocd_check
class: ArgoCDCheck
params:
url: https://argocd.homelab.local
token: ${ARGOCD_TOKEN}
# GitLab Pipeline
- module: gitlab_check
class: GitLabPipelineCheck
params:
url: https://gitlab.homelab.local
token: ${GITLAB_TOKEN}
project_ids:
- 1 # infra-config
- 2 # k8s-manifests
- 3 # ansible-playbooks
lookback_hours: 24
# 证书到期
- module: cert_check
class: CertificateCheck
params:
domains:
- blog.homelab.local
- vault.homelab.local
- rss.homelab.local
- argocd.homelab.local
- gitlab.homelab.local
warning_days: 30
critical_days: 7
告警规则
# config/alert_rules.yaml
notifier:
wecom_webhook: "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=xxx"
smtp:
host: smtp.homelab.local
port: 587
user: ops@homelab.local
password: ${SMTP_PASSWORD}
from: ops@homelab.local
recipients:
- admin@homelab.local
# 告警抑制规则:同一问题 1 小时内不重复告警
suppression:
window_minutes: 60
# critical 级别不受抑制
except_severity: [critical]
命令行入口
# run_inspection.py
#!/usr/bin/env python3
import yaml
import json
import sys
from pathlib import Path
from engine.runner import InspectionRunner
from engine.notifier import Notifier
def main():
config_dir = Path(__file__).parent / "config"
# 加载配置
with open(config_dir / "targets.yaml") as f:
targets = yaml.safe_load(f)
with open(config_dir / "alert_rules.yaml") as f:
alert_config = yaml.safe_load(f)
# 执行巡检
runner = InspectionRunner(targets)
print(f"🔧 运维巡检开始 - {len(targets['checks'])} 项检查\n")
results = runner.run_all()
# 汇总
summary = runner.summary()
print(f"\n{'='*60}")
print(f"总计: {summary['total']} 项 | ✅ {summary['healthy']} 正常 | "
f"⚠️ {summary['by_severity'].get('warning', 0)} 警告 | "
f"🔴 {summary['by_severity'].get('critical', 0)} 严重")
print(f"{'='*60}\n")
# 发送告警
notifier = Notifier(alert_config.get("notifier", {}))
notifier.send(results)
# 保存报告
report_path = Path(__file__).parent / "reports" / f"inspection-{summary['timestamp'][:10]}.json"
report_path.parent.mkdir(exist_ok=True)
with open(report_path, "w") as f:
json.dump(summary, f, indent=2, ensure_ascii=False)
# 非 0 退出码,方便 K8s CronJob 判断
sys.exit(0 if summary["issues"] == 0 else 1)
if __name__ == "__main__":
main()
容器化部署
# Dockerfile
FROM python:3.11-slim
WORKDIR /app
# 安装系统依赖
RUN apt-get update && apt-get install -y --no-install-recommends \
openssh-client && \
rm -rf /var/lib/apt/lists/*
# 安装 Python 依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
# 复制代码
COPY . .
# 入口
ENTRYPOINT ["python", "run_inspection.py"]
# requirements.txt
paramiko==3.4.0
kubernetes==29.0.0
requests==2.32.0
PyYAML==6.0.1
Jinja2==3.1.4
部署到 K8s
apiVersion: batch/v1
kind: CronJob
metadata:
name: inspection
namespace: monitoring
spec:
schedule: "0 */4 * * *" # 每 4 小时执行一次
concurrencyPolicy: Forbid
successfulJobsHistoryLimit: 3
failedJobsHistoryLimit: 10
jobTemplate:
spec:
backoffLimit: 1
template:
spec:
restartPolicy: OnFailure
containers:
- name: inspection
image: registry.homelab.local/inspection:latest
envFrom:
- secretRef:
name: inspection-secrets
volumeMounts:
- name: kube-config
mountPath: /root/.kube
readOnly: true
- name: ssh-key
mountPath: /root/.ssh
readOnly: true
volumes:
- name: kube-config
configMap:
name: kube-config
- name: ssh-key
secret:
secretName: ssh-key
紧急巡检:手动触发
# 手动触发一次巡检
kubectl create job --from=cronjob/inspection manual-inspection -n monitoring
# 查看巡检结果
kubectl logs job/manual-inspection -n monitoring
# 只检查特定项目(如只检查证书)
kubectl exec -it deployment/inspection-web -n monitoring -- \
python run_inspection.py --check cert_check
巡检报告 Web 仪表盘
用 Flask 搭一个简单的仪表盘查看历史巡检结果:
# app.py
from flask import Flask, render_template, jsonify
import json
from pathlib import Path
from datetime import datetime
app = Flask(__name__)
REPORT_DIR = Path(__file__).parent / "reports"
@app.route("/")
def dashboard():
reports = sorted(REPORT_DIR.glob("*.json"), reverse=True)[:10]
histories = []
for report in reports:
with open(report) as f:
data = json.load(f)
histories.append({
"timestamp": data["timestamp"],
"total": data["total"],
"healthy": data["healthy"],
"issues": data["issues"],
"by_severity": data["by_severity"],
})
# 最新一次
latest = histories[0] if histories else None
return render_template("report.html", histories=histories, latest=latest)
@app.route("/api/latest")
def api_latest():
reports = sorted(REPORT_DIR.glob("*.json"), reverse=True)
if reports:
with open(reports[0]) as f:
return jsonify(json.load(f))
return jsonify({"error": "no report"})
@app.route("/api/report/<date>")
def api_report(date):
report = REPORT_DIR / f"inspection-{date}.json"
if report.exists():
with open(report) as f:
return jsonify(json.load(f))
return jsonify({"error": "not found"}), 404
if __name__ == "__main__":
app.run(host="0.0.0.0", port=8080)
执行效果
$ python run_inspection.py
🔧 运维巡检开始 - 6 项检查
✅ [ OK] SSH连通性 | 5台主机 | 全部 5 台主机 SSH 可达
⚠️ [ WARNING] 磁盘空间 | 5台主机 | 1 个磁盘超过 80%
✅ [ OK] K8s集群健康 | 3 节点 | 全部 3 个节点就绪,无异常 Pod
🔴 [CRITICAL] ArgoCD同步 | 8个应用 | 2 个应用健康状态异常
✅ [ OK] GitLab Pipeline | 3个项目 | 最近 24h 内 12 个 Pipeline 全部成功
⚠️ [ WARNING] 证书到期 | 5个域名 | 1 个证书 30 天内到期
============================================================
总计: 6 项 | ✅ 3 正常 | ⚠️ 2 警告 | 🔴 1 严重
============================================================
整个巡检流程在 7 台机器 + K8s 集群上跑完大约 45 秒,比人工打开 6 个浏览器标签页要快得多。
扩展检查器
这套架构是插件化的,加一个新检查器只需要 3 步:
# engine/checks/custom_check.py
from .base import CheckResult, Severity
class CustomCheck:
"""自定义检查器:你的检查逻辑"""
def __init__(self, params: dict):
self.target = params.get("target", "")
def run(self) -> CheckResult:
# 你的检查逻辑
is_ok = True # 替换为实际检查
if is_ok:
return CheckResult(
name="自定义检查",
target=self.target,
severity=Severity.OK,
message="检查通过",
)
else:
return CheckResult(
name="自定义检查",
target=self.target,
severity=Severity.WARNING,
message="发现问题",
)
然后在 targets.yaml 里加一条配置即可:
- module: custom_check
class: CustomCheck
params:
target: "你的目标"
与现有监控的分工
| 巡检系统 |
Prometheus + AlertManager |
| 每 4 小时跑一次 |
实时监控 |
| 覆盖配置/CI/CD/证书 |
覆盖指标/性能/资源 |
| 主动拉取 |
被动接收 |
| 适合发现"慢变质" |
适合发现"突发故障" |
| 不依赖被监控端运行 |
依赖 Exporter 在线 |
两者互补:Prometheus 管"着火",巡检系统管"漏水"。着火要秒级响应,漏水要定期排查。
常见坑
坑 1:SSH 检查被防火墙挡。 巡检 Pod 在 K8s 集群里,如果目标机器的 SSH 只允许特定 IP,需要在防火墙放行 K8s 节点 IP。
坑 2:K8s API 连接超时。 如果 kubeconfig 指向外部 LB,LB 宕了就检查不了。建议巡检 Pod 部署在集群内,用 ServiceAccount 直接访问 API Server。这也是为什么上面 CronJob 示例选择用 ConfigMap 挂载 kubeconfig、而不是依赖外部入口。
坑 3:企业微信 Webhook 限流。 每分钟最多 20 条。如果检查项多,聚合后再发一条,不要每个检查器单独发。
坑 4:证书检查的域名在内网。 内网域名公网 DNS 解析不了,需要在巡检容器里配 /etc/hosts 或使用内网 DNS。
小结
运维自动化系列的 7 篇文章到这里就全部结束了。回顾一下这个系列:
| 篇 |
主题 |
核心工具 |
解决什么问题 |
| 1 |
批量管理 |
Ansible |
不用逐台 SSH |
| 2 |
基础设施即代码 |
Terraform |
资源可重建 |
| 3 |
GitOps |
ArgoCD |
声明式部署 |
| 4 |
包管理 |
Helm |
应用标准化 |
| 5 |
CI/CD |
GitLab |
代码到上线自动化 |
| 6 |
定时任务 |
K8s CronJob |
crontab 集群化 |
| 7 |
巡检系统 |
自建 |
覆盖监控盲区 |
自动化不是目的,减少重复劳动才是。如果你每天还在手动检查系统状态、手动部署、手动清理日志——是时候让机器来做了。
运维的终局是:你定规则,机器执行。你喝咖啡,系统自己跑。如果想交流更多自动化运维实践,也欢迎常来云栈社区逛逛。