
简介这是基于机器学习的分布式故障检测Python项目源码包主要面向计算机专业正在准备毕业设计、课程设计或期末大作业的学生也适合需要完整项目实战练习的学习者。项目围绕分布式环境中的异常发现与故障定位展开结合了数据预处理、特征构建和模型调用等流程并配有可视化页面用于结果展示技术路线清晰。压缩包共179个文件容量约47.95MB核心包括38个Python源码文件及编译后的pyc文件含h5/pkl模型数据、csv样本数据以及html/css/js等前端展示资源文件类型覆盖从模型训练到界面呈现的完整链路。项目已经过严格调试可直接解压运行并作为毕业设计主体节省从零搭建环境与编写代码的时间。资源目前已有121人学习下载对希望快速获得可演示、可扩展项目的读者有实际参考价值。1. 故障检测不是监控告警是让模型替你盯指标分布式系统里最头疼的问题不是“节点挂没挂”而是“节点正在变慢、变怪但还活着”。传统监控用阈值告警比如CPU超过90%就报警但如果某个节点因为磁盘IO抖动导致请求延迟从10ms涨到800msCPU却只有30%经典阈值根本抓不到。这个项目做的就是把这套判断交给机器学习先把CPU、内存、网络、日志等指标整理成特征再用随机森林、孤立森林、SVM这些分类器去识别“正常”和“故障”的边界。它面向两类人一类是做毕业设计需要完整可运行源码的学生另一类是想在真实分布式环境里落地异常检测的开发者。项目里包含了数据构造、特征提取、模型训练、检测模块和可视化基本上把一条链路都串起来了。后面我会按实际拆项目的顺序把每个环节的坑和参数逐一说明。2. 特征工程把CPU、内存和网络延迟变成模型认识的向量2.1 原始监控数据长什么样分布式故障检测的第一步不是建模而是搞清楚手头有什么数据。常见监控系统里每个节点每5秒会采集一条记录包含timestamp、node_id、cpu_usage、mem_usage、disk_io、net_throughput、latency、error_count这些字段。这个项目的数据集就是用Python模拟生成的故意在部分时间片里注入了“故障模式”比如突然的内存泄漏、网络丢包、磁盘写满。模拟数据的意义在于你知道哪个时间片是故障才能给模型标标签也才能用准确率、召回率去衡量检测效果。真实环境里你往往没有这种“标准答案”所以先用模拟数据验证流程是通用做法。2.2 滑动窗口统计特征的计算单条监控数据不足以判断故障因为瞬时抖动可能只是噪声。常见做法是构造滑动窗口把过去60秒12个采样点的数据聚合出统计量均值、标准差、最大值、最小值、变化率以及当前值相对窗口均值的偏离程度。这些特征能反映“趋势”而不是“瞬间值”对慢故障尤其有效。下面是一段典型的特征提取代码用pandas实现import pandas as pd import numpy as np def build_features(df, window_size12): df: 原始监控数据包含cpu_usage, mem_usage, latency等列 window_size: 滑动窗口大小5秒一个点12个点即60秒 返回: 拼接了窗口统计特征的新DataFrame features [] for _, group in df.groupby(node_id): # 按时间排序确保窗口顺序正确 group group.sort_values(timestamp) for col in [cpu_usage, mem_usage, disk_io, latency]: # 滚动窗口均值shift(1)避免用到当前点自身造成数据泄漏 group[f{col}_win_mean] group[col].rolling(window_size).mean().shift(1) group[f{col}_win_std] group[col].rolling(window_size).std().shift(1) group[f{col}_win_diff] group[col] - group[f{col}_win_mean] features.append(group) return pd.concat(features) # 假设df是读取的监控数据 df pd.read_csv(monitor_data.csv) feature_df build_features(df) print(feature_df.head())这段代码里最容易被忽略的是shift(1)。如果直接用rolling的均值来预测当前点特征里包含了当前点的信息训练时准确率会虚高但部署到线上就失效。win_diff表示当前值偏离窗口均值的程度这是故障检测里最核心的信号正常状态下偏离很小故障时会出现持续偏移或剧烈波动。参数window_size一般设12或24窗口太短捕捉不到慢故障太长又会把真正的异常平滑掉。如果你的监控间隔是30秒那么窗口大小就要相应调整到20个点以上才够覆盖10分钟。2.3 标签构造怎么定义“故障”特征有了标签怎么来真实系统里故障不是非黑即白但训练分类器必须有确定标签。这个项目采用的办法是在模拟阶段设定故障注入区间比如第1000到1200个采样点为“内存泄漏故障”第2000到2200个为“网络抖动故障”。落在这个区间的样本标记为label1其余为0。注意一个细节故障不会瞬间发生通常有一个“劣化”过程所以代码里会对故障标签做前向膨胀把故障开始前的20个点也标为1。这样模型才能学到“异常前的征兆”而不是只学到故障爆发后的样子。def make_label(df, fault_windows): df: 带时间戳的监控数据 fault_windows: 列表每个元素是 (node_id, start_time, end_time, fault_type) 返回: 加上label列的DataFrame df[label] 0 for node, start, end, _ in fault_windows: mask (df[node_id] node) (df[timestamp] start) (df[timestamp] end) df.loc[mask, label] 1 # 前向膨胀把故障开始前20个点也标为1 pre_start start - 20 * 5 # 假设5秒一个点20个点即100秒 pre_mask (df[node_id] node) (df[timestamp] pre_start) (df[timestamp] start) df.loc[pre_mask, label] 1 return df标签膨胀比例需要控制好太大模型会过度关注早期信号误报增加太小则模型学不到“临近故障”的状态。实践里我一般先试20个点然后看训练曲线再调整。这个项目源码里写的是PRE_FAULT_POINTS 20直接改这个常量就能重训。3. 模型选型与训练随机森林、孤立森林和SVM的取舍3.1 为什么不用阈值告警有人会问CPU升高就告警这不是很简单吗问题在于分布式系统的故障往往是复合型的。比如某个节点网络延迟升高但CPU和内存都正常另一个节点磁盘IO高但延迟正常。单指标阈值需要为每个指标手工调参上百个节点的生产环境根本维护不过来。机器学习模型可以自动组合指标间的关联——比如“网络延迟高 CPU低 磁盘写等待高”这组模式比单一阈值可靠得多。而且随机森林这类模型还能输出特征重要度告诉你哪些特征对判断故障贡献最大这比拍脑袋定阈值有说服力。3.2 三模型对比和关键参数这个项目里训练了三个模型随机森林、孤立森林和SVM。它们的定位不同随机森林RandomForest监督学习需要标签适合故障类型已知、能打标的数据。抗过拟合强特征重要性可直接查看。孤立森林IsolationForest无监督不需要标签适合标签缺失的探索阶段。它基于“异常点更容易被隔离”的思想但会把没见过的正常状态也当异常误报偏高。SVMRBF核小样本分类效果好但对特征缩放敏感分布式场景下特征维度高时训练速度慢。下表是项目里实际用的参数直接抄来可跑模型参数数值说明RandomForestn_estimators200树的数量越大越稳但慢RandomForestmax_depth12限制单棵树深度防过拟合RandomForestmin_samples_split5内部节点再划分所需最小样本数RandomForestclass_weightbalanced正负样本不平衡时自动加权IsolationForestcontamination0.05期望的异常比例需先估计IsolationForestn_estimators200孤立树数量SVMC1.0越大越易过拟合越小容忍误分类SVMgammascale自动按特征方差缩放SVMkernelrbf非线性分类常用3.3 训练脚本与模型保存下面是训练和保存的核心代码可以看到数据切分和评估也一并处理了from sklearn.ensemble import RandomForestClassifier, IsolationForest from sklearn.svm import SVC from sklearn.model_selection import train_test_split from sklearn.metrics import classification_report, confusion_matrix import joblib def train_models(feature_df): feature_cols [c for c in feature_df.columns if c not in [node_id, timestamp, label]] X feature_df[feature_cols].fillna(0) # 窗口前几行会产生NaN填0 y feature_df[label] # 按时间切分不随机打乱避免用未来数据训过去 split_idx int(len(X) * 0.8) X_train, X_test X.iloc[:split_idx], X.iloc[split_idx:] y_train, y_test y.iloc[:split_idx], y.iloc[split_idx:] # 随机森林 rf RandomForestClassifier( n_estimators200, max_depth12, min_samples_split5, class_weightbalanced, n_jobs-1 ) rf.fit(X_train, y_train) print(RandomForest test report:) print(classification_report(y_test, rf.predict(X_test))) # 孤立森林需要先训练异常检测器再用它生成新特征这里直接训练并输出 iso IsolationForest(contamination0.05, n_estimators200, random_state42) iso.fit(X_train) # 无监督只喂特征 # 孤立森林输出1为正常-1为异常转为0/1标签 iso_pred (iso.predict(X_test) -1).astype(int) print(IsolationForest confusion matrix:) print(confusion_matrix(y_test, iso_pred)) # 保存随机森林模型 joblib.dump(rf, rf_fault_detector.pkl) joblib.dump(iso, iso_fault_detector.pkl) return rf, iso这段代码有两点需要注意一是切分方式按时间顺序切分而不是随机切分。故障检测本质是时间序列预测随机切分会把前后样本混在一起容易让模型“偷看”未来信息导致评估结果虚高。二是孤立森林不训练标签所以它和随机森林的预测语义不同。实战中我一般先跑孤立森林做一轮无监督筛查把认为是异常的时间段挑出来人工看一眼确认是真实故障后再打标签训练随机森林。SVM的训练代码类似但要先做特征标准化否则RBF核会直接失效from sklearn.preprocessing import StandardScaler scaler StandardScaler() X_train_s scaler.fit_transform(X_train) X_test_s scaler.transform(X_test) svm SVC(C1.0, kernelrbf, gammascale, class_weightbalanced) svm.fit(X_train_s, y_train) print(SVM report:) print(classification_report(y_test, svm.predict(X_test_s)))标准化只对SVM这类基于距离的模型必要对树模型不是必须。项目里把这些封装在train_pipeline.py里改参数后重新运行即可。不要忽略fillna(0)滚动窗口前几行是NaN如果不处理sklearn会直接报错。4. 分布式检测模块把模型部署到每个节点和中心节点4.1 整体架构Agent采集、中心聚合训练好的模型要真正跑在分布式系统里不能只在离线脚本里做实验。这个项目给出了一套可运行的检测模块结构分两层Agent层部署在每个被监控节点上负责采集指标、用已训练的模型做本地快速判断并上报结果。中心层接收所有Agent的报告做综合判定比如超过半数节点认为异常就触发告警。这种设计的优点是单节点故障时检测不依赖中心Agent自己就能发现本地异常中心层则能捕捉到跨节点的关联故障比如一个机柜温度过高导致多个节点同时异常。分布式的故障检测必须考虑这一点单机异常可能是自身问题多个节点同时异常往往是环境问题。4.2 Agent端实时检测代码Agent端一般用定时任务实现每5秒拉取一次本机指标构造当前特征然后调用模型预测。这里给出一个简化但可运行的Agent伪代码import psutil import time import joblib import pandas as pd class FaultAgent: def __init__(self, node_id, model_path, window_size12): self.node_id node_id self.model joblib.load(model_path) self.window_size window_size self.history [] # 保存最近window_size个原始指标 def collect_metrics(self): # 用psutil拉取CPU、内存、磁盘IO等指标 cpu psutil.cpu_percent(interval1) mem psutil.virtual_memory().percent disk psutil.disk_io_counters().read_bytes / 1024 / 1024 # MB latency self._measure_latency() # 自定义方法测到某个依赖服务的延迟 return {cpu_usage: cpu, mem_usage: mem, disk_io: disk, latency: latency} def _measure_latency(self): # 简化实现用socket测到中心节点的响应时间 import socket, time start time.time() try: socket.create_connection((10.0.0.1, 9999), timeout1).close() return (time.time() - start) * 1000 # ms except Exception: return 1000 # 超时视为1000ms def predict(self): m self.collect_metrics() self.history.append(m) if len(self.history) self.window_size: return False # 窗口未满暂不判断 self.history self.history[-self.window_size:] # 构造特征需要和训练时保持一致 df pd.DataFrame(self.history) features {} for col in [cpu_usage, mem_usage, disk_io, latency]: features[f{col}_win_mean] df[col].mean() features[f{col}_win_std] df[col].std() features[f{col}_win_diff] df[col].iloc[-1] - df[col].mean() pred self.model.predict(pd.DataFrame([features]).fillna(0))[0] return bool(pred) agent FaultAgent(node-01, rf_fault_detector.pkl) if agent.predict(): print(f[{time.time()}] node-01 is abnormal!)这里_measure_latency模拟了网络延迟的测量生产环境里可能是从连接池或调用链追踪数据里拿。注意Agent端的history不能无限增长要始终保留最近window_size个点否则内存会泄漏。窗口长度必须和训练时一致这里的window_size12对应训练代码里的rolling(12)。4.3 中心节点投票与告警中心节点逻辑更简单每个Agent定时上报自己的预测结果异常/正常中心维护一个状态表统计最近5分钟或者最近N个周期每个节点异常的次数。如果某个节点异常次数超过阈值或者超过一半节点异常就触发告警并写入日志。from collections import defaultdict, deque class CenterNode: def __init__(self, alert_threshold3, degradation_ratio0.5): # 记录每个节点最近的异常状态最多保留10个周期 self.status defaultdict(lambda: deque(maxlen10)) self.alert_threshold alert_threshold self.degradation_ratio degradation_ratio def on_report(self, node_id, is_abnormal): self.status[node_id].append(is_abnormal) total_nodes len(self.status) abnormal_nodes sum( 1 for node in self.status if sum(self.status[node]) self.alert_threshold ) if abnormal_nodes / total_nodes self.degradation_ratio: print(f[ALERT] {abnormal_nodes}/{total_nodes} nodes abnormal!) return True return False这个实现里alert_threshold是重点参数。设成3表示一个节点连续3个周期15秒都报异常才认定为故障这样可以过滤掉偶发抖动。degradation_ratio设成0.5表示超过一半节点异常时触发集群级告警这能捕捉到网络分区或机房断电这类大面积故障。实际部署时告警信息可以接钉钉、企业微信或者发邮件源码里用print代替方便跑通流程。5. 调参和踩坑从“能跑”到“真的能发现问题”5.1 样本不平衡怎么处理分布式系统正常运行时间远远多于故障时间故障样本可能只占1%。直接训练随机森林会把所有样本都预测成正常准确率99%但毫无意义。除了用class_weightbalanced还可以在训练前下采样正常样本。项目里给了一个下采样函数建议在故障样本量充足时使用from sklearn.utils import resample def balance_dataset(X, y): X_normal, X_abnormal X[y 0], X[y 1] # 让正常样本和异常样本数量一致 X_normal_res resample(X_normal, replaceFalse, n_sampleslen(X_abnormal), random_state42) X_bal pd.concat([X_normal_res, X_abnormal]) y_bal pd.concat([pd.Series([0]*len(X_normal_res)), pd.Series([1]*len(X_abnormal))]) return X_bal, y_bal注意replaceFalse如果故障样本太少则改为replaceTrue用有放回采样。下采样会丢失大量正常数据模型可能丢失正常状态多样性导致误报增加。另一种更稳的做法是不动数据用class_weight调权重我一般先试class_weight效果不满意再下采样。5.2 误报率压不下去怎么办误报在故障检测里比漏报更让人头疼因为频繁告警会让人麻木最后真正故障时没人响应。压误报可以从三个方向入手提高孤立森林的contamination参数不行正确的是降低它。contamination是模型期望的异常比例设高了会误把正常边缘状态划为异常。调整判决阈值。随机森林默认用0.5作为正负类分界但你可以输出预测概率然后选更高的阈值比如0.7才报异常。这牺牲召回率换精确率适合告警场景from sklearn.metrics import precision_recall_curve probs rf.predict_proba(X_test)[:, 1] precision, recall, thresholds precision_recall_curve(y_test, probs) # 找到精确率不低于0.9的最高召回率对应的阈值 valid [(p, r, t) for p, r, t in zip(precision, recall, thresholds) if p 0.9] best max(valid, keylambda x: x[1]) if valid else None if best: threshold best[2] print(fSelected threshold: {threshold:.3f})检查特征时间对齐问题。如果Agent上报的数据有延迟或丢失窗口内会混入不完整数据导致预测值异常。中心节点应该丢弃超过2秒延迟的报告而不是直接使用。5.3 用混淆矩阵验证检测效果最后验证模型不能只看准确率。故障检测是典型的类别不平衡问题我建议每次训练后都打印混淆矩阵并确认下面四个数from sklearn.metrics import confusion_matrix cm confusion_matrix(y_test, rf.predict(X_test)) tn, fp, fn, tp cm.ravel() print(fTP{tp} (真实故障、检出了)) print(fFN{fn} (漏报最危险)) print(fFP{fp} (误报会烦死人)) print(fTN{tn} (正常判断正常))FN是漏报意味着故障没被发现后果最严重FP是误报会消耗运维精力。实际业务里如果漏报和误报的代价不同应该用ROC曲线下面积或PR曲线来选模型而不是准确率。这个项目源码里evaluate.py会输出这三张图你跑一遍就能直观看到随机森林在F1分数上通常优于SVM因为SVM处理高维连续特征时如果某些特征分布不是高斯型RBF核的表现会很不稳定。孤立森林作为无监督基线召回率往往偏高但精确率低适合做第一道粗筛再用随机森林细排。本文还有配套的精品资源点击获取