
简介这份资源是面向Java开发者与分布式系统运维人员的机器学习故障诊断项目源码适合具备一定Java基础、希望将机器学习方法落地到系统监控与异常排查场景的中级学习者。项目以Java为主要实现语言围绕分布式环境下的故障识别与诊断流程组织代码可用于课程设计、毕业设计或工程实践参考。压缩包共33个文件其中27个java源文件承载核心算法与业务逻辑5个xml与1个yml文件负责依赖配置与项目参数管理整体约23KB结构紧凑、便于快速通读。目前已有358人学习下载。读者可从中获取一套可运行的故障诊断工程骨架理解特征提取、模型训练与诊断结果输出之间的代码衔接方式并借助pom.xml等配置快速导入IDE调试为后续扩展数据集或替换算法模型提供清晰的改造入口。1. 从一份 Java 源码说起分布式系统故障诊断到底难在哪线上八台订单服务节点凌晨两点开始每隔十几分钟就有一台响应时间从 80ms 飙到 2s日志里没有 ERRORCPU、内存、GC 指标都正常重启后能好一阵子又复发。这种场景做过分布式系统的人大概率都遇到过它最折磨人的地方不是故障本身而是你手里有一堆监控数据却不知道该看哪个。基于 Java 机器学习的分布式系统故障诊断系统要解决的就是这件事把散落在各个节点上的指标、日志、调用链数据收集起来用机器学习模型自动判断当前系统处于什么状态、哪个环节出了问题。这份源码适合两类人看一类是想把机器学习真正落到运维场景的 Java 后端另一类是做课程设计或毕业设计需要一套完整可跑系统的同学。它不是一个调几个 API 就完事的 demo里面涉及数据采集、特征工程、模型训练、在线推理、告警联动这一整条链路每一环都有它自己的坑。2. 拆开这套系统数据从哪来、特征怎么造、模型怎么选2.1 分布式故障诊断的数据采集层怎么搭分布式系统的数据源大致分三类。第一类是指标数据比如每个节点的 CPU 使用率、内存占用、磁盘 IO、网络吞吐、JVM 堆使用、GC 次数和耗时这些通常通过 Micrometer 或 Prometheus 客户端暴露成 HTTP 端点再由采集器定时拉取。第二类是日志数据包括应用日志和系统日志关键是从中提取出异常堆栈、超时记录、连接池耗尽这类事件。第三类是链路数据也就是一次请求经过了哪些服务、每个环节耗时多少OpenTelemetry 或 SkyWalking 都能提供。这套源码里我一般会把采集层做成一个独立的 Spring Boot 模块用定时任务拉取指标用 Logback Appender 把日志异步写到消息队列链路数据则通过 Agent 上报。采集频率是个需要认真对待的参数指标类数据 15 秒一次比较合适太密会给被监控系统增加负担太疏会漏掉短时抖动。日志和链路数据是事件驱动的来一条收一条但要加缓冲队列防止突发流量打爆采集端。// 指标采集任务每 15 秒执行一次 Component public class MetricCollector { private final MeterRegistry registry; private final KafkaTemplateString, String kafkaTemplate; // 采集频率通过配置文件注入默认 15000 毫秒 Scheduled(fixedRateString ${collector.metric.interval:15000}) public void collect() { // 遍历所有已注册的指标打包成 JSON 发送到 Kafka registry.getMeters().forEach(meter - { MetricSnapshot snapshot new MetricSnapshot(); snapshot.setName(meter.getId().getName()); snapshot.setTags(meter.getId().getTags().toString()); snapshot.setTimestamp(System.currentTimeMillis()); // 只采集可量化的数值型指标 if (meter instanceof Gauge) { snapshot.setValue(((Gauge) meter).value()); } else if (meter instanceof Counter) { snapshot.setValue(((Counter) meter).count()); } kafkaTemplate.send(metric-topic, JSON.toJSONString(snapshot)); }); } }这段代码的关键点在于fixedRateString用了占位符方便不同环境调整采集频率。MetricSnapshot里保留了标签信息因为后续做故障定位时需要按节点、按服务维度聚合。发送到 Kafka 而不是直接写库是为了削峰和解耦采集端不需要关心下游是实时消费还是离线训练。2.2 特征工程把原始指标变成模型能吃的输入原始指标直接丢给模型效果通常很差因为故障信号往往藏在变化趋势和指标之间的关联里。我一般会做三层特征。第一层是统计特征对每个指标在滑动窗口内算均值、方差、最大值、最小值、P95 分位值窗口大小取 1 分钟、5 分钟、15 分钟三档。第二层是差分特征算一阶差分和二阶差分用来捕捉突增突降。第三层是关联特征比如 CPU 使用率和请求延迟的比值、GC 耗时占 CPU 时间的比例这类交叉特征对区分“真故障”和“正常波动”特别有用。// 滑动窗口统计特征计算 public class FeatureExtractor { // 窗口大小1分钟、5分钟、15分钟单位毫秒 private static final long[] WINDOWS {60_000L, 300_000L, 900_000L}; public MapString, Double extract(ListMetricPoint points, long now) { MapString, Double features new HashMap(); for (long window : WINDOWS) { // 过滤出窗口内的数据点 ListDouble values points.stream() .filter(p - p.getTimestamp() now - window) .map(MetricPoint::getValue) .collect(Collectors.toList()); if (values.isEmpty()) continue; String prefix w (window / 60_000) _; features.put(prefix mean, values.stream().mapToDouble(d - d).average().orElse(0)); features.put(prefix std, std(values)); features.put(prefix max, values.stream().mapToDouble(d - d).max().orElse(0)); features.put(prefix p95, percentile(values, 95)); // 一阶差分均值反映变化速度 features.put(prefix diff_mean, diffMean(values)); } return features; } }窗口大小的选择直接影响模型灵敏度。1 分钟窗口能抓住秒级抖动但噪声也大15 分钟窗口平滑但可能漏掉短故障。三个窗口一起用让模型自己学哪个窗口的特征更重要。diffMean这个特征在诊断“缓慢劣化型”故障时特别管用比如内存泄漏导致堆使用率持续上升均值可能还在正常范围但差分均值已经明显偏正。2.3 模型选型为什么随机森林和 XGBoost 是首选故障诊断本质上是一个分类问题给定当前特征判断系统处于正常、CPU 瓶颈、内存瓶颈、网络延迟、磁盘 IO 瓶颈中的哪一种。可选算法很多但落地时我优先考虑随机森林和 XGBoost原因有三个。第一它们对特征尺度不敏感不需要做归一化省去很多预处理麻烦。第二训练速度快几万条样本几分钟就能出模型方便迭代。第三特征重要性可以直接输出运维人员能看懂“这次报警主要是因为哪个指标异常”这对建立信任很关键。深度学习不是不能用LSTM 做时序异常检测效果确实好但它需要更多数据、更长训练时间而且推理时延比树模型高一个数量级。在线诊断场景下从指标异常到发出告警最好控制在秒级树模型更容易满足这个要求。如果你们的数据量足够大、故障模式特别复杂可以先用树模型上线再慢慢积累数据迭代到深度模型。# 训练脚本核心部分用 XGBoost 做多分类 import xgboost as xgb from sklearn.model_selection import train_test_split from sklearn.metrics import classification_report # X 是特征矩阵y 是故障标签 0-正常 1-CPU 2-内存 3-网络 4-磁盘 X_train, X_test, y_train, y_test train_test_split(X, y, test_size0.2, stratifyy) model xgb.XGBClassifier( n_estimators200, # 树的数量太少欠拟合太多容易过拟合 max_depth6, # 树深度故障诊断一般 4-8 够用 learning_rate0.1, # 学习率配合 n_estimators 调整 subsample0.8, # 行采样比例防止过拟合 colsample_bytree0.8, # 列采样比例 objectivemulti:softmax, num_class5 ) model.fit(X_train, y_train, eval_set[(X_test, y_test)], early_stopping_rounds20) print(classification_report(y_test, model.predict(X_test)))early_stopping_rounds20是个实用技巧验证集损失连续 20 轮不下降就停避免无效训练。stratifyy保证训练集和测试集的类别比例一致故障样本通常远少于正常样本不做分层可能导致测试集里故障样本太少评估结果不可信。3. 从源码到跑通本地环境搭建与最小可运行链路3.1 环境依赖与版本选择这套系统依赖的组件不算少但都是主流选择。JDK 用 11 或 17 都行17 的 ZGC 在采集端高吞吐场景下更有优势。Spring Boot 2.7.x 是稳定选择3.x 对 JDK 版本要求更高如果团队还没升级可以先用 2.7。消息队列用 Kafka单机跑一个 broker 就够验证。存储层分两块原始指标和日志存 Elasticsearch 方便检索特征和模型元数据存 MySQL。Python 环境用 3.8 以上XGBoost 和 scikit-learn 装最新稳定版即可。本地验证时不需要把所有组件都搭起来可以先用内存队列替代 Kafka用 H2 替代 MySQL把核心链路跑通再逐步替换成真实组件。我一般会准备一个docker-compose.yml把 Kafka、Elasticsearch、MySQL 一键拉起来省去手动安装的麻烦。# docker-compose.yml 核心片段 version: 3.8 services: kafka: image: bitnami/kafka:3.4 ports: - 9092:9092 environment: - KAFKA_CFG_NODE_ID0 - KAFKA_CFG_PROCESS_ROLEScontroller,broker - KAFKA_CFG_LISTENERSPLAINTEXT://:9092,CONTROLLER://:9093 elasticsearch: image: docker.elastic.co/elasticsearch/elasticsearch:8.7.0 ports: - 9200:9200 environment: - discovery.typesingle-node - xpack.security.enabledfalse mysql: image: mysql:8.0 ports: - 3306:3306 environment: - MYSQL_ROOT_PASSWORDroot123 - MYSQL_DATABASEfault_diagnosisKafka 用 KRaft 模式不需要 ZooKeeper单节点就能跑。Elasticsearch 关掉安全认证方便本地调试生产环境一定要开。MySQL 建一个fault_diagnosis库表结构在源码的sql目录下导入即可。3.2 启动顺序与配置要点组件启动有先后依赖。先起 Kafka、Elasticsearch、MySQL再起采集端 Spring Boot 应用最后起诊断服务。采集端启动时会检查 Kafka 连接连不上会重试但诊断服务如果连不上 MySQL 会直接启动失败所以数据库必须先就绪。配置文件里几个关键参数需要根据环境调整。collector.metric.interval控制采集频率本地验证可以调到 5000 毫秒加快数据产生。kafka.bootstrap-servers指向 Kafka 地址。diagnosis.model.path指定模型文件路径首次启动时如果模型文件不存在诊断服务会跳过加载等训练完成后手动触发 reload。# 启动顺序 docker-compose up -d kafka elasticsearch mysql # 等待约 30 秒让组件完全就绪 sleep 30 # 启动采集端 java -jar collector/target/collector-1.0.0.jar --spring.profiles.activelocal # 启动诊断服务 java -jar diagnosis/target/diagnosis-1.0.0.jar --spring.profiles.activelocal启动后访问诊断服务的/actuator/health确认状态再访问/api/diagnosis/status看模型是否加载成功。如果模型未加载调用/api/diagnosis/train触发一次训练训练数据来自 Elasticsearch 里已采集的指标。3.3 用模拟数据验证整条链路真实故障不好复现验证阶段我一般会写一个模拟器人为制造 CPU 飙高、内存泄漏、网络延迟等场景。模拟器通过 JMX 或直接调用系统接口修改指标值让采集端能抓到异常数据。// 故障模拟器制造 CPU 高负载场景 Component public class FaultSimulator { Value(${simulator.enabled:false}) private boolean enabled; // 模拟 CPU 飙高启动多个忙循环线程 public void simulateCpuSpike(int threads, long durationMs) { if (!enabled) return; for (int i 0; i threads; i) { new Thread(() - { long end System.currentTimeMillis() durationMs; while (System.currentTimeMillis() end) { // 空转消耗 CPU Math.sqrt(Math.random()); } }).start(); } } }simulator.enabled默认关闭只在验证环境打开。threads控制并发线程数一般设为 CPU 核数加一就能把使用率拉到 90% 以上。durationMs控制持续时间建议至少 60 秒让采集端能收集到足够多的异常样本。跑完模拟后观察诊断服务是否在预期时间内发出 CPU 瓶颈告警如果没发检查特征计算窗口是否覆盖了异常时间段。4. 避坑指南这套系统落地时最容易翻车的五个地方4.1 告警风暴模型频繁切换状态导致误报现象是诊断服务在几分钟内反复发出和撤销告警运维人员被大量通知淹没。原因是模型对每个采集点独立判断指标在阈值附近波动时分类结果会在正常和异常之间跳变。解决办法是加状态保持机制连续 N 个采集周期都判定为异常才真正触发告警N 一般取 3 到 5。同时加冷却时间同一类告警触发后 10 分钟内不再重复发。4.2 特征穿越训练时用了未来数据导致线上效果暴跌现象是离线评估准确率 95% 以上上线后准确率不到 60%。原因是特征计算时用了当前时刻之后的数据比如算 5 分钟窗口均值时把未来 5 分钟的点也算进去了。解决办法是特征计算严格按时间戳过滤只取timestamp now的数据点并且训练集和测试集按时间切分而不是随机切分避免同一时间段的数据同时出现在训练和测试中。4.3 样本不均衡正常样本太多导致模型只会说“正常”现象是模型把所有输入都判为正常故障召回率接近零。原因是正常样本通常是故障样本的几十倍甚至上百倍。解决办法有三个方向一是对故障样本过采样二是对正常样本欠采样三是在损失函数里给故障类别更高权重。XGBoost 的scale_pos_weight参数可以调整正负样本权重多分类场景下可以用sample_weight给每个样本赋权。4.4 模型更新断层新故障模式出现后模型无法识别现象是系统上线新版本后出现了一种以前没见过的故障模型把它判为正常。原因是模型只见过训练数据里出现过的故障类型对未知模式没有识别能力。解决办法是加一个异常检测兜底层用孤立森林或自编码器判断当前特征是否偏离历史分布偏离度超过阈值就触发“未知异常”告警同时把样本记录下来供后续标注和重新训练。4.5 采集端成为瓶颈监控系统本身拖垮被监控系统现象是接入诊断系统后业务服务响应变慢排查发现是采集端占用了过多资源。原因是采集频率太高或采集的数据量太大。解决办法是采集端独立部署不和业务服务混部采集频率根据指标重要性分级核心指标 15 秒一次非核心指标 60 秒一次日志采集用异步 Appender避免阻塞业务线程。5. 让模型持续可用的几个进阶技巧模型上线只是开始真正决定这套系统能不能长期用的是后续的迭代机制。我一般会做三件事。第一件是建立反馈闭环每次告警发出后让运维人员标记“确认故障”或“误报”这些标记数据定期回流到训练集。第二件是监控模型自身的健康度统计每天的预测分布如果正常样本占比突然从 95% 掉到 80%说明要么系统真出了大问题要么模型漂移了两种情况都需要人工介入。第三件是定期做影子模式验证新模型上线前先让它跑在后台做预测但不发告警对比新旧模型的判断差异确认新模型没有明显退化再切换。# 模型漂移检测对比近期预测分布与训练分布 import numpy as np from scipy.stats import entropy def detect_drift(recent_predictions, train_distribution, threshold0.1): # 统计近期预测的类别分布 recent_dist np.bincount(recent_predictions, minlengthlen(train_distribution)) recent_dist recent_dist / recent_dist.sum() # 计算 KL 散度值越大说明分布差异越大 kl_div entropy(recent_dist, train_distribution) if kl_div threshold: return True, kl_div # 需要重新训练 return False, kl_divthreshold设 0.1 是个经验值低于这个值说明分布基本一致高于这个值就要警惕。KL 散度对零概率敏感计算前要给每个类别加一个极小值做平滑。这个检测每天跑一次结果记录到监控面板上趋势持续上升就安排重新训练。还有一个容易被忽略的点是模型版本管理。每次训练产出的模型文件要带版本号和训练数据时间范围诊断服务加载模型时记录当前版本告警信息里带上模型版本方便回溯问题时确认是哪个模型做的判断。我吃过这个亏有一次线上误报率突然升高排查半天才发现是有人手动替换了模型文件但没记录最后只能靠文件修改时间倒推。特征重要性的定期审查也很有价值。如果某个特征的重要性突然从第十名跳到第一名往往意味着数据采集出了问题比如某个指标的单位变了或者采集脚本改了。这种变化不一定是坏事但必须搞清楚原因再决定是否继续用。希望帮到你。本文还有配套的精品资源点击获取