专栏简介: 本专栏将从单机巡检起步,逐步演进到多机批量巡检、异常告警推送、可视化Dashboard,最终构建一套完整的企业级自动化运维平台。每一篇都是生产环境实战沉淀,拒绝拒绝Hello World式教学。
📌 本篇导读
如果沿用上一篇企业级Python自动化运维实战(二):单机升级到多机巡检系统,我们需要在每台服务器上都执行一遍,进行三次ssh链接,效率极低。本篇我们将迎来V2.1版本的重大升级,将单机脚本演进为多机并发巡检系统。我们将引入SSH密钥认证、多线程并发、结构化数据解析以及完善的日志系统,打造真正符合企业级标准的自动化运维工具。
| 认证方式 | 密码/账号明文有风险 | SSH Key 免密认证 |
| 执行效率 | 单台服务器逐个检查 | 多线程并发执行 |
| 容错能力 | SSH失败程序直接退出 | 异常容错,单机失败不影响整体 |
| 数据格式 | 采集结果是字符串 | 结构化数据(JSON/Dict) |
| 日志系统 | 无日志记录 | 标准化日志系统 |
| 报告质量 | 简单文本 | 标准化HTML报告 |
一、 安全基石:SSH Key 免密认证
在生产环境中,明文管理服务器密码是极大的安全隐患。企业级运维的标准做法是使用 SSH Key(公钥/私钥) 进行认证。
1.1 生成 SSH 秘钥
在我们的巡检服务器(控制机)上执行以下命令,一路回车即可生成密钥对:
作者:饭要加加加
如需代码压缩包请关注作者,私信获取
ssh-keygen -t rsa -b 4096
执行后,~/.ssh/ 目录下会生成 id_rsa(私钥)和 id_rsa.pub(公钥)。
1.2 分发公钥
将公钥复制到目标服务器上,实现免密登录:
ssh-copy-id root@10.207.132.221
1.3 配置管理
修改 config/server.yaml,将认证方式从密码改为密钥文件路径:
servers:
– name: master
host: 10.207.132.221
username: root
key_file: C:\\Users\\Administrator\\.ssh\\id_rsa
1.4 升级 SSH 模块
修改 ssh/client.py,使用 paramiko 加载私钥进行连接。
import paramiko
def ssh_execute(server, command):
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
# 加载私钥
private_key = paramiko.RSAKey.from_private_key_file(server["key_file"])
client.connect(
hostname=server["host"],
username=server["username"],
pkey=private_key,
timeout=5
)
stdin, stdout, stderr = client.exec_command(command)
result = stdout.read().decode()
client.close()
return result
二、 健壮性升级:增加 SSH 异常处理
在V1版本中,如果某台服务器IP错误或宕机,程序会直接崩溃退出。在企业环境中,单机故障绝不能影响整体巡检任务的执行。
我们需要对 ssh_execute 函数进行“防御性编程”改造,捕获所有异常并返回结构化结果:
import paramiko
def ssh_execute(server, command):
try:
client = paramiko.SSHClient()
client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
key = paramiko.RSAKey.from_private_key_file(server["key_file"])
client.connect(
hostname=server["host"],
username=server["username"],
pkey=key,
timeout=5
)
stdin, stdout, stderr = client.exec_command(command)
result = stdout.read().decode()
client.close()
# 成功返回 success 状态
return {
"status": "success",
"data": result
}
except Exception as e:
# 失败返回 failed 状态及错误信息,程序不崩溃
return {
"status": "failed",
"error": str(e)
}
💡 核心知识点:Paramiko 常用功能
| SSH登录 | SSHClient | connect() | 建立连接,传入IP、端口、用户名、密钥等。 |
| 执行命令 | SSHClient | exec_command() | 执行命令,返回 stdin, stdout, stderr 三个通道。 |
| 文件上传 | SFTPClient | put() | 上传文件 (本地路径, 远程路径)。 |
| 文件下载 | SFTPClient | get() | 下载文件 (远程路径, 本地路径)。 |
| 密钥认证 | RSAKey | from_private_key_file() | 加载私钥文件,配合 connect(pkey=) 使用。 |
三、 性能飞跃:重构数据采集与多线程并发
3.1 采集策略重构
之前的做法是“采集一个指标建立一次SSH连接”,这在多机巡检中是致命的性能瓶颈。 优化方案:创建一个新的入口 collectors/server.py,一次SSH连接,批量执行所有命令,然后在本地进行解析。
# collectors/server.py
from ssh.client import ssh_execute
from parser.cpu import parse_cpu
from parser.memory import parse_memory
from parser.disk import parse_disk
def inspect_server(server):
# 一次性执行多个命令
command = """
echo CPU
top -bn1 | grep Cpu
echo MEMORY
free -m
echo DISK
df -h /
"""
result = ssh_execute(server, command)
if result["status"] != "success":
return {
"name": server["name"],
"status": "failed",
"error": result["data"]
}
raw_data = result["data"]
# 本地解析数据
return {
"name": server["name"],
"cpu": parse_cpu(raw_data),
"memory": parse_memory(raw_data),
"disk": parse_disk(raw_data),
"status": "success"
}
3.2 数据结构化:Parser 模块
SSH 返回的是字符串,程序无法直接判断大小。我们需要使用正则表达式(re 模块)将其转换为数字。
创建 parser/ 目录:
parser/cpu.py
import re
def parse_cpu(data):
result = re.search(r'(\\d+\\.\\d+)\\s*id', data)
if result:
idle = float(result.group(1))
return round(100 – idle, 2)
return 0
parser/disk.py
import re
def parse_disk(data):
result = re.search(r'(\\d+)%\\s+/', data)
if result:
return float(result.group(1))
return 0
parser/memory.py
import re
def parse_memory(data):
result = re.search(r'Mem:\\s+(\\d+)\\s+(\\d+)', data)
if result:
total = float(result.group(1))
used = float(result.group(2))
if total > 0:
return round(used / total * 100, 2)
return 0
💡 核心知识点:Python 正则 re 模块
| 全局搜索(首个) | re.search() | 扫描整个字符串,返回第一个匹配结果(最常用)。 |
| 全局搜索(全部) | re.findall() | 返回列表,包含所有匹配结果。 |
| 从头匹配 | re.match() | 仅从字符串开头匹配,不常用。 |
| 字符串替换 | re.sub() | 替换匹配的字符串。 |
| 编译对象 | re.compile() | 预编译正则,提高重复使用时的性能。 |
3.3 多线程并发执行
当面对 100 台服务器时,串行执行可能需要 500 秒(假设每台5秒),这是不可接受的。我们使用 concurrent.futures 实现多线程并发。
创建 core/task.py:
from concurrent.futures import ThreadPoolExecutor, as_completed
def run_parallel(tasks, max_workers=10):
"""
tasks: [(函数, 参数), …]
max_workers: 最大线程数
"""
results = []
with ThreadPoolExecutor(max_workers=max_workers) as executor:
future_list = []
for func, args in tasks:
future = executor.submit(func, args)
future_list.append(future)
for future in as_completed(future_list):
try:
result = future.result()
results.append(result)
except Exception as e:
results.append({"status": "failed", "error": str(e)})
return results
💡 核心知识点:ThreadPoolExecutor
| 创建线程池 | ThreadPoolExecutor(max_workers=N) | 创建包含 N 个线程的池子。 |
| 提交任务 | submit(fn, *args) | 提交任务,立即返回 Future 对象。 |
| 关闭线程池 | shutdown(wait=True) | 等待任务完成并关闭(通常配合 with 使用)。 |
关于 Future 对象:
- future.result(): 阻塞等待任务完成并获取返回值。
- future.done(): 判断任务是否完成。
- future.exception(): 获取任务执行时的异常。
四、 可观测性:增加日志系统
企业级项目必须有日志。我们将使用 Python 内置的 logging 模块记录巡检过程。
在 main.py 中配置日志:
import logging
logging.basicConfig(
filename="logs/app.log",
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s"
)
# 使用示例
logging.info("开始巡检 web01")
logging.error("web01 连接失败")
日志输出示例: 2026-07-29 15:00 INFO 开始巡检web01 2026-07-29 15:01 ERROR web02 timeout
五、 最终整合:Main 入口与报告调整
最后,我们将所有模块串联起来。同时,调整 analyzer/check.py 以适应新的数据结构(包含 name, status 等字段)。
analyzer/check.py
CPU_THRESHOLD = 80
MEMORY_THRESHOLD = 80
DISK_THRESHOLD = 90
def check(data):
alarms = []
# 先判断采集是否成功
if data.get("status") != "success":
alarms.append("服务器数据采集失败")
return alarms
# 指标检查
if data.get("cpu", 0) > CPU_THRESHOLD:
alarms.append("CPU使用率过高")
if data.get("memory", 0) > MEMORY_THRESHOLD:
alarms.append("内存使用率过高")
if data.get("disk", 0) > DISK_THRESHOLD:
alarms.append("磁盘空间不足")
return alarms
main.py
import os
import sys
import logging
from config.ssh_config import load_servers
from report.report import create_report
from core.task import run_parallel
from analyzer.check import check
from collectors.server import inspect_server
def main():
logging.basicConfig(
filename="logs/app.log",
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s"
)
servers = load_servers()
tasks = []
# 构建任务列表
for server in servers:
tasks.append((inspect_server, server))
# 并发执行
results = run_parallel(tasks)
# 处理结果
for data in results:
alarms = check(data)
data["alarms"] = alarms
# 生成报告
create_report(results)
if __name__ == "__main__":
main()
六.本篇总结与下篇预告
✅ 本篇收获
至此,我们的巡检系统完成了从“单机脚本”到“企业级工具”的华丽转身:
🚀 下篇预告:数据库与Web监控平台改造 系统已经跑起来了,但报告还是太简陋?告警只能看日志? 在下一篇文章中,我们将进行数据持久化到数据库,制作精美的 HTML 可视化报告,并接入 企业微信/钉钉机器人,实现真正的实时告警推送!
🌟 创作不易,如果这篇文章对你有帮助,欢迎点赞 👍、收藏 ⭐、关注 🔔 三连支持!你的支持是我持续更新的最大动力。有任何问题或建议,欢迎在评论区留言交流,我会一一回复!
网硕互联帮助中心


评论前必须登录!
注册