在一台连不上外网的机器上把 4TB 原始运行日志变成一份能跑起来的模型这件事听起来像是纯算法问题实际干下来你会发现八成时间花在集群、内存和 ETL 上。这篇文章讲的就是这类军用大数据 Spark 机器学习场景下的完整落地链路数据不出内网、硬件不是最新、依赖全靠离线搬运怎么把 Spark 集群搭起来怎么把日志洗成特征矩阵怎么用 MLlib 训出模型以及那些文档里不会写、只有真跑过一遍才会遇到的坑。适合正在做大数据毕设、准备大数据面试、或者刚接手一个封闭环境下机器学习任务的同学看也适合已经会写单机 scikit-learn、但对分布式训练还没概念的人。题目的关键词是 Spark、机器学习、大数据三个词凑在一起最容易出现的误解就是Spark 有 MLlib所以模型也能像单机一样随手跑。真做过一版就知道分布式机器学习真正的难点从来不在算法本身而在资源怎么分、数据怎么切、shuffle 怎么压、结果怎么复现。军工、能源、通信这类对数据外流零容忍的行业天然没有云、没有托管服务、没有在线调参平台所有东西都得自己在内网里动手搭。下面按实际动手顺序展开。1. 封闭内网里选型为什么是 Spark 而不是别的1.1 数据体量和迭代次数决定了技术栈的下限先算一笔账这笔账决定了你到底需不需要 Spark。假设原始遥测日志每天 30GB文本格式压缩比大约 4 到 5 倍一年下来压缩后大约 2TB 出头。如果做滑窗特征窗口 5 秒、步长 1 秒特征维度做到 300 到 500 维最终训练矩阵的规模很容易到几亿行。单机 128GB 内存的服务器用 pandas 读一份 20GB 的 Parquet 就已经在悬崖边上了再做一次 groupby rolling 基本必崩。这里有个经常被忽略的点Spark 的价值不只是数据放得下。更关键的是它把特征工程也纳入了同一套执行引擎。你在单机上习惯的读数据、洗一遍、算特征、训练四步在分布式环境下最贵的成本是每一步之间的落盘和再读取。Spark 的 DataFrame 加上 MLlib 的 Pipeline能让清洗、特征、训练串成一条 DAG中间结果尽量留在内存里只在 shuffle 时落盘。这才是它在这类场景里不可替代的地方。所以选型判断标准很简单如果一份特征矩阵单机能装进内存并且训练时间在可接受范围内就别上 Spark。用 DuckDB 或者 Polars 加单机 scikit-learn开发和调试效率高一个数量级。只有当数据必须分片、或者特征生成本身就是个几小时的批处理任务时Spark 才真正开始划算。1.2 版本选择2.4 还是 3.x这是个真问题内网环境最尴尬的一点是你往往不能自由选择版本因为要迁就已有的 Hadoop 集群、操作系统和 JDK。下表是过去几年我实际遇到过的几种组合以及最终的选择现有环境推荐 Spark 版本理由主要代价Hadoop 2.7 / JDK 8Spark 2.4.x官方二进制包直接兼容改动最小无 AQE无 Arrow 加速Pandas UDF 受限Hadoop 3.x / JDK 8Spark 3.2 - 3.3AQE 可大幅缓解倾斜SQL 优化成熟需要重新对齐依赖、重测现有脚本无 Hadoop纯本地盘Spark 3.3 Standalone部署最轻不依赖 YARN资源调度能力弱适合固定作业国产化操作系统与 JDK 版本严格对齐多为 JDK 8 环境部分 native 库缺失需自行编译我对这一段的建议是优先选 Spark 3.x只要 JDK 允许。原因很实在——Adaptive Query ExecutionAQE在 3.x 里默认开启后能自动合并小分区、自动处理部分倾斜 join这一项就能省掉大量手工调参的时间。2.4 时代为了处理一个倾斜 key你可能要写加盐逻辑改三四处代码3.x 里很多时候改个spark.sql.adaptive.enabledtrue加两个阈值就过去了。1.3 部署模式Standalone、YARN、Kubernetes 怎么挑内网项目里我几乎没见过用 Kubernetes 的原因很简单运维成本高而作业数量少、形态固定。三种模式的取舍大致是这样Standalone装起来最快集群规模在 5 到 50 台之间最舒服。缺点是资源是静态划分的多作业并发时容易互相抢。适合每天晚上跑一批固定作业这种节奏。YARN如果公司/单位已经有 Hadoop 集群直接复用是最省事的。队列、资源池、权限都是现成的。缺点是调试链路长一个 OOM 你要翻 YARN 日志、Spark 日志、executor 日志三层。Kubernetes动态资源、镜像化环境确实优雅但在完全离线的环境里光是维护一个私有镜像仓库和网络插件就能耗掉半个月。我接手这个项目时集群是 6 台机器1 主 5 从用的 Standalone。原因很实际没有现成的 YARN也不值得为了十几个批处理作业去引入一整套调度体系。2. 把 Standalone 集群从零拉起来2.1 节点角色划分master 不要跑 executor很多人第一次搭 Spark 会把 master 节点也当成 worker 用理由是机器空着浪费。在小规模测试时看不出问题一旦作业量上来master 既要负责资源协商、又要跑 executor很容易出现 driver 心跳超时。我的做法是master 节点只跑 Master 进程SPARK_WORKER_CORES设成 0 或者干脆不在这台机器上启动 Worker。worker 节点上每台的 executor 数量按核数 ÷ 每 executor 核数来定通常每个 executor 给 4 到 5 核比较合适。磁盘方面Spark 的 shuffle 是写本地盘的所以 worker 节点的本地盘 IOPS 比容量更重要。我遇到过一台机器 shuffle 阶段特别慢最后查出来是那块盘在做 RAID 重建。2.2 离线搬运依赖清单比安装步骤更重要内网部署最耗时的不是装是凑齐依赖。我的经验是先在能联网的机器上把所有东西列一张清单一次性下齐避免来回搬优盘。清单大致包括JDK 安装包版本要和 Spark 编译时的目标一致Spark 3.x 大多用 JDK 8 或 11Spark 二进制包注意选带 Hadoop 版本的那个包名比如spark-3.3.2-bin-hadoop3.tgzPython 环境如果集群上要用 PySpark各节点的 Python 版本必须一致。我用 conda-pack 把整个环境打包后分发比逐台 pip install 靠谱得多。如果要用到额外算法库或 native 依赖把 wheel 文件也一起带上离线pip install --no-index --find-links。提示先在每台机器上用sha256sum校验安装包完整性。从优盘多次拷贝大文件损坏是真实发生过的。2.3 配置文件三个文件决定集群能不能稳Spark Standalone 的核心配置就三个文件但每一条都有讲究# conf/spark-env.sh export JAVA_HOME/opt/jdk1.8.0_xxx export SPARK_MASTER_HOST10.x.x.10 export SPARK_MASTER_PORT7077 export SPARK_MASTER_WEBUI_PORT8080 export SPARK_WORKER_CORES32 export SPARK_WORKER_MEMORY200g export SPARK_WORKER_DIR/data/spark/work export SPARK_LOCAL_DIRS/data/spark/tmp# conf/workers 10.x.x.11 10.x.x.12 10.x.x.13 10.x.x.14 10.x.x.15SPARK_WORKER_MEMORY不要填成物理内存的全量。留下至少 10% 到 15% 给操作系统、页缓存和其他进程。SPARK_LOCAL_DIRS一定要指到大盘上默认的/tmp在很多系统上只有几十 GBshuffle 一放大就写满作业直接失败报错信息还很难指向真正的原因。spark-defaults.conf里我通常先落这几条基线spark.serializerorg.apache.spark.serializer.KryoSerializer spark.sql.adaptive.enabledtrue spark.sql.shuffle.partitions800 spark.sql.parquet.compression.codecsnappy spark.locality.wait3s spark.network.timeout600s spark.executor.heartbeatInterval60sspark.sql.shuffle.partitions的默认值 200 在数据量大时明显不够但也不是越大越好。经验值是让每个分区处理 100 到 200MB 的数据。你可以先按这个估算跑几次看 stage 的 task 耗时分布再调整。2.4 启动验证失败时报错指向哪里启动流程本身很标准master 上跑sbin/start-master.sh然后sbin/start-workers.sh。验证按这个顺序走能快速定位问题jps看进程。master 上应该有Masterworker 上应该有Worker。打开http://master:8080看 worker 是否都注册上了Alive Workers 数量对不对。跑官方示例验证计算链路bin/spark-submit --class org.apache.spark.examples.SparkPi --master spark://master:7077 examples/jars/spark-examples_*.jar 100常见的启动失败只有几类我按出现频率排一下hostname 无法解析/etc/hosts没写全或者用了短名和 FQDN 混着写、SSH 免密没配好导致start-workers.sh拉不起来远程进程、端口被占用、worker 配置的内存超过实际可用内存导致注册失败。最后这一类最容易误判——Web UI 上看到 worker 在线但一提交作业就报 not enough resources实际原因是SPARK_WORKER_MEMORY设得比物理可用内存还大Worker 注册时报告的资源是假的。3. 内存与资源Spark 调优里最容易翻车的一环3.1 executor 的内存到底花在哪这个问题在大数据面试里被问烂了但真正让人栽跟头的是细节。一个 executor 从操作系统角度看占用的内存是spark.executor.memoryspark.executor.memoryOverheadspark.executor.pyspark.memoryPySpark 场景而spark.executor.memory内部又分成三块保留内存默认 300MB、用户内存、Spark 内存。Spark 内存再按spark.memory.fraction默认 0.6划分为执行内存和存储内存其中存储部分由spark.memory.storageFraction默认 0.5控制。举个实际配置例子。一台 256GB 内存、48 核的机器我会这样配每个 executor 给 5 核spark.executor.cores5每台起 8 个 executorspark.executor.memory24gspark.executor.memoryOverhead4g合计 8 × 28g 224g剩下的 32g 留给系统和页缓存但 8 × 5 40 核留了 8 核给系统合理注意这里memoryOverhead给到了 4g比例比默认的 10% 高。原因是用了 PySparkPython 进程的内存是算在 overhead 里的默认值很容易触发Container killed by YARN for exceeding memory limits这类报错Standalone 下表现是 executor 被 Worker 杀掉。3.2 数据倾斜识别比解决更难倾斜的典型表现是一个 stage 里 99% 的 task 三分钟跑完剩下几个跑了四十分钟或者 Spark UI 上看到某个 task 的 shuffle read 是其他 task 的几十倍再或者 executor 的 GC 时间占比一直高于 10%。识别出来之后处理手段按代价从低到高排让 AQE 自动处理。Spark 3.x 下开启spark.sql.adaptive.skewJoin.enabledtrue并调整spark.sql.adaptive.skewJoin.skewedPartitionFactor很多中等程度的倾斜就自动拆了。两阶段聚合。先对 key 加随机前缀做一次局部聚合再去掉前缀做全局聚合。这个手法对groupBy类操作特别有效代码改动不大from pyspark.sql import functions as F # 第一阶段加盐局部聚合 salted df.withColumn(salt, (F.rand() * 16).cast(int)) \ .withColumn(salted_key, F.concat_ws(_, device_id, salt)) partial salted.groupBy(salted_key).agg(F.sum(value).alias(partial_sum)) # 第二阶段去盐全局聚合 final partial.withColumn(device_id, F.split(salted_key, _)[0]) \ .groupBy(device_id).agg(F.sum(partial_sum).alias(total))广播小表。如果是 join 导致的倾斜且其中一张表很小默认阈值 10MB可调spark.sql.autoBroadcastJoinThreshold直接广播掉。拆开单独处理。如果倾斜来自少数几个异常 key比如某个设备 ID 因为采集程序 bug 写了几亿条记录最务实的做法是把这批 key 单独捞出来处理剩下的走正常流程。工程上难看但有效。3.3 序列化与 shuffle 参数几行配置换半小时Kryo 序列化是必开的而且要注册自定义类否则 Kryo 会退化成写全类名。虽然 Spark 内置了常见类的注册但你的自定义 case class 或者 JavaBean 如果不注册序列化开销反而可能比 Java 序列化还大。shuffle 相关的几个参数我在大批量作业里会这样调spark.shuffle.file.buffer1m spark.reducer.maxSizeInFlight96m spark.shuffle.compresstrue spark.shuffle.spill.compresstrue spark.memory.fraction0.7 spark.memory.storageFraction0.4spark.memory.fraction从默认 0.6 提到 0.7是因为这类批处理作业基本不用缓存 DataFrame把内存更多让给 shuffle 和聚合操作相应地降低storageFraction避免存储内存长期占着不放。3.4 报错到原因的对账表报错信息真实原因处理方向java.lang.OutOfMemoryError: GC overhead limit exceeded执行内存不足对象频繁晋升老年代提高 executor 内存减少每分区数据量优化倾斜java.lang.OutOfMemoryError: Java heap space单分区数据过大或 collect 到 driver检查是否有collect()增大分区数FetchFailedExceptionshuffle 落盘文件丢失或节点压力大检查SPARK_LOCAL_DIRS磁盘减小分区粒度ExecutorLostFailureexecutor 被杀通常伴随内存超限提高 memoryOverhead检查 Worker 日志Task not serializable闭包捕获了不可序列化的对象把外部对象声明为 transient或用广播变量Container killed by YARN...超过容器内存上限同时调 memory 和 memoryOverhead4. 从原始日志到特征矩阵ETL 才是主战场4.1 数据源与分区策略小文件是隐形杀手原始数据通常是几种形态混合文本日志、CSV、从数据库导出的 AVRO还有一种最麻烦的——固件直接吐出的二进制遥测帧需要按协议手动解析。不管哪种读完第一步永远是按日期分区重写成 Parquet。分区数量有个经验值让每个分区文件在 128MB 到 256MB 之间。太少会导致单个 task 处理量过大太多会产生海量小文件让 NameNode 和后续读取都变慢。写的时候用repartition而不是coalesce来控制文件数因为coalesce只能减少分区且不保证均匀df.repartition(F.col(dt), F.ceil(F.rand() * 8).cast(int)) \ .write.mode(overwrite).partitionBy(dt) \ .parquet(/warehouse/telemetry_parquet)4.2 时间对齐与清洗最容易埋下数据泄漏的地方清洗环节我固定做这几件事时间戳统一到毫秒并转成 UTC设备本地时间和采集时间混用是常见错误来源、重复记录按(device_id, ts)去重保留最后一条、缺失值按设备和指标分组做前向填充、超量程的读数直接置空而不是替换成均值。这里有个非常容易被忽略的点滑窗特征会造成标签泄漏。假如你的标签是未来 5 秒是否异常特征窗口如果跨到了标签时间窗内离线评估的 AUC 会漂亮得离谱上线就崩。我在第一次做的时候就这么踩过AUC 0.98上线之后掉到 0.6。解决办法是在生成特征时就显式加一条时间约束特征窗口的右边界必须早于标签窗口的左边界至少一个缓冲期。4.3 特征工程的三个层次我把特征分成三类分开设计、分开验证统计类特征滑窗内的均值、标准差、峰度、过零率、极差。用window窗口函数实现注意窗口的 partitionBy 一定要带上设备 ID否则会跨设备计算。from pyspark.sql.window import Window w Window.partitionBy(device_id).orderBy(F.col(ts).cast(long)) \ .rangeBetween(-300, -1) # 过去 300 秒 feat df.withColumn(v_mean, F.avg(v).over(w)) \ .withColumn(v_std, F.stddev(v).over(w))频域特征对振动、电流这类周期性信号取一段定长序列做 FFT抽出主频、频谱能量占比几个特征。Spark 本身没有原生 FFT通常的做法是把序列按groupBy collect_list聚成数组后用 Pandas UDF 批量处理。这里一定用向量化的 Pandas UDF 而不是逐行的udf性能差十倍以上。类别与交叉特征StringIndexer加OneHotEncoder是标准组合。有个细节——StringIndexer在训练集上拟合完成后如果线上出现没见过的类别默认会报错。生产环境要设handleInvalidkeep或者提前准备一份完整的字典并固定映射。4.4 Pipeline 组装与落盘特征组装用VectorAssembler注意所有输入列必须是数值类型布尔列要先转成double。我一般会把整个特征流程写成 Pipeline 并保存下来因为训练和打分必须用完全一致的处理逻辑靠人工保证最终一定会不一致。from pyspark.ml import Pipeline from pyspark.ml.feature import VectorAssembler, StandardScaler assembler VectorAssembler( inputCols[c for c in feat.columns if c.startswith(v_)], outputColraw_features, handleInvalidskip ) scaler StandardScaler(inputColraw_features, outputColfeatures, withMeanTrue, withStdTrue) pipeline Pipeline(stages[assembler, scaler]) feature_pipeline pipeline.fit(train_df) feature_pipeline.write().overwrite().save(/models/feature_pipeline_v1)特征矩阵落 Parquet不要落 CSV。Parquet 有 schema、有压缩、有列裁剪读取速度差好几倍而且不会因为某个字段里有逗号就整体错位——这个坑我用 CSV 踩过一次排查了两个小时。5. 用 MLlib 建模算法、切分与评估5.1 算法选型的现实考虑在这类场景下我选算法的第一标准不是准确率上限而是可解释性和稳定性。原因很直接模型出问题的时候你需要向非技术的同事解释为什么这台设备被判成异常一个特征重要性能画出来的模型沟通成本比黑盒低太多。算法适用场景优势注意点逻辑回归二分类基线、需要概率输出可解释、训练快需手动处理非线性特征要归一化GBDT结构化特征的主力效果好、给特征重要性树多时训练慢需控制深度和迭代数随机森林特征维度高、噪声多抗过拟合、并行度好模型体积大预测延迟高KMeans无标签的设备分组快、直观k 值需实验对量纲敏感ALS设备与工况的关联推荐稀疏数据友好冷启动问题LDA日志文本主题抽取无监督、可解释主题数需调中文需先分词我的常规路径是先用逻辑回归跑一个基线确认特征和标签的关系是合理的再上 GBDT 看能提升多少如果 GBDT 相对 LR 提升在 3 个点以内直接上 LR省下来的推理时间和维护成本更值。5.2 数据切分时间序列千万别随机切这是我在这个项目里最想强调的一条。传感器数据是有时间相关性的随机切分会让训练集和验证集里存在几乎相同的样本评估结果虚高。正确做法是按时间切比如前 8 个月训练第 9 个月验证第 10 个月测试。更进一步如果做了滑窗特征相邻样本之间本身就有重叠切分点附近要有隔离带。我通常会在切分点前后各丢掉一个窗口长度虽然损失一些样本但能显著降低评估偏差。5.3 交叉验证的代价要提前算MLlib 的CrossValidator用起来很方便但代价是参数组合数 × fold 数 × 单次训练时间。一个 GBDT 单次训练 20 分钟5 折 × 8 组参数 40 次训练 13 个小时。在内网环境里这还算好的因为资源独享如果是共享集群可能要跑两天。我的做法是分两步先用单次留出验证在少量候选参数上粗筛选出两三个方向再对这几个方向做 3 折交叉验证精调。这样总时间能压到原来的三分之一。from pyspark.ml.tuning import ParamGridBuilder, CrossValidator from pyspark.ml.evaluation import BinaryClassificationEvaluator from pyspark.ml.classification import GBTClassifier gbt GBTClassifier(featuresColfeatures, labelCollabel, maxIter80) grid ParamGridBuilder() \ .addGrid(gbt.maxDepth, [4, 6]) \ .addGrid(gbt.stepSize, [0.05, 0.1]) \ .build() evaluator BinaryClassificationEvaluator( labelCollabel, rawPredictionColrawPrediction, metricNameareaUnderROC) cv CrossValidator(estimatorgbt, estimatorParamMapsgrid, evaluatorevaluator, numFolds3, parallelism4) cv_model cv.fit(train_df)5.4 评估指标与阈值AUC 高不等于能用不平衡数据是常态异常样本可能只占 0.5%。这种情况下 AUC 只能说明排序能力不能说明业务效果。真正要盯的是 PR 曲线和不同阈值下的召回率、误报率。我会做一张阈值对账表把业务侧的容忍度翻译成阈值如果业务能接受每千条里 5 条误报那就在验证集上找误报率不超过 0.5% 的前提下召回率最高的那个阈值。这件事必须在离线阶段做完上线之后再调阈值风险高得多。5.5 模型持久化与跨环境加载model.save()存的是 Spark ML 的格式加载时如果 Spark 版本不一致或者自定义的 UDF 类路径找不到会直接失败。我的经验是训练和推理用同一个集群的同一个 Spark 版本模型目录里额外放一份特征处理流程的代码快照。纯手工但出问题的时候能救命。离线打分时注意VectorAssembler的列顺序必须和训练时完全一致。这一点 Pipeline 已经帮你保证了前提是你加载的是保存过的那条 Pipeline而不是重新拼一条看起来一样的。6. 跑通之后调度、监控和那些反复出现的问题6.1 幂等与断点重跑批处理作业一定要幂等。做法是先把结果写到临时目录全部成功后用原子操作替换正式目录失败就直接删临时目录。这样重跑不会产生半份数据也不会污染下游。另外一个实用技巧是按日期做断点重跑。作业入口接受一个日期范围参数每次只处理还没成功的日期。我曾经因为一个夜跑任务连续三天失败导致补数据要手工跑六次后来把断点逻辑加上之后一条命令就补完了。6.2 该监控什么Spark 作业的监控不需要太复杂但有几项必须有作业总时长和每个 stage 的耗时对比历史基线超过 30% 就告警。shuffle 落盘总量突然翻倍通常意味着分区数不对或者数据分布变了。executor 的 GC 时间占比长期高于 10% 说明内存配置需要重新算。输入数据量如果某天数据量掉了一半很可能是上游采集出了问题而不是业务变好了。我在实际运维中遇到过最隐蔽的一次是某个设备的数据从某天开始全部变成同一个值。作业没报错、指标没异常、耗时也正常但训练出来的模型对这类设备完全失效。后来加了一条每个分区的某个关键字段的唯一值数量监控才提前发现这类数据漂移。6.3 无外网环境下的实验记录没有在线实验平台实验记录只能靠自己。我用的是最土但最可靠的办法每个实验一个目录目录里放四样东西——提交脚本、特征处理代码的 Git commit hash、参数配置文件、评估结果 JSON。同时在本地搭一个轻量的实验跟踪工具内网部署的 MLflow 完全够用把所有指标汇总到一张表里。这件事的意义在三个月后才体现出来当有人问上个月那版模型和现在这版差别在哪你能在五分钟内说清楚而不是重新翻聊天记录。注意实验目录不要放在 Spark 的工作目录下面。SPARK_LOCAL_DIRS会被清理实验记录丢了就真的找不回来了。6.4 一次典型的连续超时排查最后说一个具体的排查过程因为它的思路比结论更有价值。现象某个夜跑作业连续三天超时从原来的 40 分钟涨到 3 小时以上。第一步看 Spark UI发现卡在一个groupBy的 stage 上task 耗时分布严重不均最慢的几个 task 是普通 task 的 20 倍。第二步看 shuffle read慢 task 读到了 2GB 以上其他只有 30MB 左右典型的倾斜。第三步去查为什么这个 key 会这么大把该 key 的数据捞出来一看是某个设备的device_id因为采集程序的一个字段拼接 bug把多个设备的记录挂到了同一个 ID 下面数据量是正常设备的几百倍。第四步修 bug 是一方面但当天的数据还得处理于是临时用两阶段聚合把这次作业跑通同时加了一条监控——单个 key 的数据量超过阈值就告警。整个过程里最有价值的一步其实是第三步。很多人看到倾斜就直接加盐、调参数问题暂时缓解了但上游的 bug 一直还在后面还会以别的形式回来。6.5 一些零碎的实操建议关于 Spark 内存我个人的体会是不要迷信网上的黄金比例。每套集群、每种数据形态的最优配置都不一样能做的只是给一个合理起点然后靠跑出来的 GC 时间、spill 大小去迭代。我一般会固定做三组对比实验内存不变换分区数、分区数不变换内存、换序列化方式跑完基本就能定下来。关于特征工程最重要的是先把评估链路搭好再去做特征。我见过太多人花两周做了一百多个特征最后发现评估指标本身有问题白干。正确的顺序是用最简单的几个特征跑通全流程确认切分、评估、复现都没问题再开始加特征。关于工具内网环境里最值得投入的两件事是把集群初始化和依赖安装写成一键脚本省下的时间远超写脚本的时间以及把常用的 ETL 和调参过程封装成可复用的函数避免每个项目重新写一遍。这两件事我在第二个项目里做完之后从拿到新数据到跑出第一版模型的周期从两周压缩到了三天。
