4.4 Pipeline与模型持久化:流程打包与模型存活


4.4 Pipeline 与模型持久化:流程打包与模型存活

本节摘要:Pipeline 把特征链与算法封装成一条可整体 fit、整体保存、整体复放的流水线;持久化让训练产物离开作业进程长期存活。本节拆解 fit 与 transform 在 Pipeline 上的两阶段执行语义,走过保存加载的完整闭环,并交代训练评估的防泄漏细节。

上一节把特征铸成了向量,本节解决两个工程收尾问题:一堆零散的 Transformer 和 Model 怎么打包不散架,训练好的模型怎么在作业结束后继续服务。这两个问题的答案在引擎视角下尤其清晰——Pipeline 就是一张"可序列化的 DataFrame 血统模板"。

Pipeline:把零件串成可保存的图纸

from pyspark.ml import Pipeline from pyspark.ml.feature import StringIndexer, VectorAssembler from pyspark.ml.classification import LogisticRegression stages = [ StringIndexer(inputCol="city", outputCol="city_idx"), VectorAssembler(inputCols=["city_idx", "age", "income"], outputCol="features"), LogisticRegression(featuresCol="features", labelCol="churn", maxIter=20) ] pipe = Pipeline(stages=stages) model = pipe.fit(train_df) # 两阶段语义的入口

fit 的执行过程按阶段逐个推进:遇到 Transformer 直接登记;遇到 Estimator 就用当前数据训练出 Model 再把 Model 装进流水线。所以一次 pipe.fit 里 StringIndexer 扫了全量数据统计类别频次,LogisticRegression 又扫了全量做迭代——fit 的成本是各阶段训练成本之和,并不存在什么"整体优化"。

transform 则完全另一副面孔:整条链全是 Transformer,逐样本窄依赖流水线,一次穿堂过。训练贵、打分便宜的分野在这里落到实锤。

pred = model.transform(test_df) # 训练 12 分钟的链,打分 40 秒 pred.select("churn", "prediction", "probability").show(5)

持久化:模型离开作业活下来

# 保存:整个 PipelineModel 一次性落盘,特征规则与模型权重同进同出 model.write().overwrite().save("hdfs://nn:9000/models/churn_v3") # 下个作业加载:换语言、换集群、隔半年,复放结果一致 from pyspark.ml.pipeline import PipelineModel loaded = PipelineModel.load("hdfs://nn:9000/models/churn_v3") loaded.transform(new_df)

持久化的关键设计是整体性:类别编号规则、向量拼装顺序、模型权重打包在一个目录里,加载时原样复放。反例是把 StringIndexer 的编号映射存在数据库、模型权重另存文件,两边版本错位时打分结果悄无声息地错——这类事故在特征链越长时越隐蔽。巡检习惯:保存目录带版本号,加载时记录版本与训练数据的指纹(行数、特征统计快照),线上打分异常时先对指纹。

训练评估:别让测试集漏进流水线

train_df, test_df = df.randomSplit([0.8, 0.2], seed=7) # 危险写法:先对全量 fit StringIndexer 再切分——类别统计见过测试集 # 正确写法:先切分,Pipeline 只 fit 训练侧 model = pipe.fit(train_df) evals = model.transform(test_df) from pyspark.ml.evaluation import BinaryClassificationEvaluator auc = BinaryClassificationEvaluator( labelCol="churn", rawPredictionCol="rawPrediction") \ .evaluate(evals) print(f"test AUC = {auc:.4f}") # 输出示例:test AUC = 0.8312

泄漏点常在流水线的第一个 Estimator:任何"先看全量数据再切分"的顺序都让测试信息渗入训练。规则一句话——切分永远在最前,fit 只碰训练侧。

环节 对象 引擎行为 常见事故
fit Pipeline 逐阶段训练,多次全量扫描 全量 fit 造成泄漏
transform PipelineModel 窄依赖流水线 上线特征与训练特征错位
保存 PipelineModel 元数据与权重整体落盘 拆开存导致版本错位
加载 任意作业 反序列化重建对象图 跨版本兼容性破裂

巡检案例:一次线上打分漂移的追查

背景:流失预测模型上线三个月,AUC 从 0.83 滑到 0.71,排障时模型权重与代码都没改动记录。操作:调出保存目录的指纹快照,与当前线上特征统计比对。结果:特征统计里 income 列的均值漂移 40%——上游数仓在两个月前改了该字段的口径(分转元),而 Pipeline 里的 LogisticRegression 权重还是按旧口径训练的。解读:Pipeline 保住了"特征规则的形状",保不住"输入数据的世界";数据漂移要靠指纹巡检发现,模型要有定期重训的节奏。处置:特征层加统计告警,触发阈值即重跑 fit 并按新版本号落盘。变式:若漂移来自类别新增(新城市上线),StringIndexer 的 unseen 处理策略(报错或兜底桶)要在训练时就定好,别等线上炸了再补。

⚠️ 常见坑:Pipeline 保存目录直接覆盖写。回滚无门是运维大忌,永远带版本号新写目录,旧版保留至少一个重训周期。

💡 关键直觉:把 PipelineModel 当成"冻结的执行计划"——它不是配置文件,是一段可复放的分布式程序图纸,输入一变它的输出语义就变,指纹巡检因此和版本管理同样重要。

本节要点回顾

  • fit 逐阶段、transform 一条线:训练成本是各 Estimator 之和,打分是窄依赖穿堂过
  • 整体持久化:特征规则与权重同进同出,拆开存等于埋版本错位的雷
  • 切分在最前:先切数据再 fit,堵住统计泄漏的主通道
  • 指纹巡检:输入世界的漂移模型自己不会报,统计快照对账是唯一手段
  • 版本化保存:永不覆盖,可回滚是模型运维的底线

MLlib 到此收束。下一章进入带结构的数据极端——图计算,看关系数据如何降维成三张 RDD 在同一台引擎上跑。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U