Spark MLlib从 RDD 到 DataFrame 的机器学习实践指南
一、MLlib 是什么?它到底能做什么?
从官方定义来看:
MLlib 是 Spark 的机器学习库(Machine Learning Library,ML),
目标是让实用的机器学习在分布式环境中既可扩展又易于使用。
从功能角度看,MLlib 提供的是一套完整的机器学习工具箱,涵盖了从算法到工程化的方方面面:
1. 机器学习算法(ML Algorithms)
MLlib 内置了一系列常见的机器学习算法,包括但不限于:
- 分类:如 Logistic Regression、决策树、随机森林、梯度提升树、Naive Bayes 等
- 回归:线性回归、决策树回归、随机森林回归、GBT 回归等
- 聚类:KMeans、Gaussian Mixture、PowerIterationClustering 等
- 协同过滤:基于 ALS 的推荐算法等
这让你可以直接在 Spark 上完成常规机器学习任务,无需再额外引入分布式训练框架。
2. 特征工程(Featurization)
工业级的机器学习,往往特征工程比模型更重要。MLlib 提供了丰富的特征处理组件,例如:
- 特征提取:词袋模型、TF-IDF、n-gram、Word2Vec 等
- 特征变换:标准化、归一化、离散化、分箱(QuantileDiscretizer)等
- 降维与特征选择:如 PCA 等
这些组件可以被组合进流水线(Pipeline),形成端到端的特征处理流程。
3. Pipelines(机器学习流水线)
Pipelines 是 Spark MLlib 的核心概念之一,它把机器学习过程抽象为一条流水线:
原始数据 → 特征处理(若干 Transformer)→ 模型训练(Estimator)→ 模型输出(Model)
通过 Pipeline,你可以:
- 把特征工程 + 模型训练组合成一个整体
- 统一保存 / 加载整个 Pipeline
- 在训练环境、生产环境之间复用同一套特征流程
4. 模型持久化(Persistence)
MLlib 提供了模型、算法、Pipeline 的保存与加载能力:
- 训练好的模型可以持久化到磁盘 / HDFS
- 在生产环境中直接加载并用于预测
- 支持 Pipeline 级别的保存,做到「一次训练,多次部署」
5. 辅助工具(Utilities)
为了支持上述能力,MLlib 还提供了:
- 线性代数:向量、矩阵、BLAS 操作等
- 统计工具:描述性统计、分布抽样等
- 数据处理:与 DataFrame、Dataset 深度集成,方便把 SQL / ETL 与 ML 结合使用
二、API 变迁:RDD-based API 进入维护模式,DataFrame 成为主力
最早接触 MLlib 的同学,可能都是通过 spark.mllib 这一套 RDD API 开始的,例如:
import org.apache.spark.mllib.regression.LabeledPoint
import org.apache.spark.mllib.classification.LogisticRegressionWithLBFGS
但是,从 Spark 2.0 开始,官方做了一个重要的战略调整:
MLlib 的 RDD-based API 已经进入“维护模式(Maintenance Mode)”。
推荐使用的主力 API,是基于 DataFrame 的spark.ml。
具体含义是:
- RDD API 仍然可以用,但只会修 Bug,不再新增功能
- 所有新的特性基本都会加在 DataFrame-based API 里
- Spark 2.x 的目标是:让 DataFrame API 与 RDD API 在功能上逐步实现 特性对齐
因此,对于一个要新做 Spark ML 的项目,正确的姿势是:优先使用 spark.ml(DataFrame API)。
三、为什么要切到 DataFrame API?RDD 不香了吗?
很多人会有疑问:RDD 这么经典的抽象,为什么要换成 DataFrame?
官方给出了几个非常关键的原因。
1. DataFrame 自身的优势
DataFrame 带来的不只是「表格化」的数据结构,更是整套 Spark 的「统一抽象」:
-
统一的数据结构:DataFrame 在 SQL、Dataset、数据源(CSV / JSON / Parquet / Hive 等)之间自然流转
-
统一的语言层 API:Python / Scala / Java / R 都有相对一致的使用体验
-
底层优化更强:
- Catalyst 优化器:对执行计划做分析与优化
- Tungsten:更高效的内存管理和代码生成
这直接带来的好处是:机器学习不再孤立,而是完全融入 Spark 生态。
2. 统一的 ML API:Pipeline + Transformer + Estimator
DataFrame API 通过三个核心概念抽象整个 ML 过程:
- Estimator:可训练的算法(如 LogisticRegression)
- Transformer:转换器,输入 DataFrame,输出新的 DataFrame(如 StringIndexer、StandardScaler)
- Pipeline:把多个 Estimator / Transformer 串起来的一条流水线
这种统一的抽象有几个明显好处:
- 不同算法、不同特征处理在使用方式上高度一致
- 更适合复杂业务场景(多步特征处理、多模型组合等)
- 更易于持久化、部署和复用
3. 更适合构建可落地的 ML Pipeline
DataFrame + Pipeline 的组合,极大降低了从「实验代码」到「生产系统」之间的距离:
- 在 Notebook 中调试好的 Pipeline,可以直接保存用于生产
- 在流式处理和批处理场景中,可以复用统一的特征逻辑
- 运维层面:模型版本、Pipeline 版本都更清晰可管理
四、“Spark ML” 是什么?和 MLlib 什么关系?
官方明确说明:
“Spark ML” 不是一个正式的产品名字,而是大家对 MLlib 中 DataFrame-based API 的一个习惯性称呼。
主要原因有两个:
- DataFrame API 的 Scala 包名是:
org.apache.spark.ml - 早期文档中常用「Spark ML Pipelines」这个说法来强调 Pipeline 概念
因此我们可以简单理解成:
-
MLlib:整个机器学习库的统称
- 包含:
spark.mllib(RDD API)和spark.ml(DataFrame API)
- 包含:
-
Spark ML:通常指
spark.ml这一套 DataFrame 风格的机器学习 API
五、MLlib 有没有被弃用(Deprecated)?
答案是:没有。
- RDD-based API 虽然进入了维护模式,但没有被标记为 deprecated
- DataFrame-based API 是目前的主推方向
- MLlib 整体依然是 Spark 中的核心组件之一
对于我们来说,你可以放心:
- 老项目里用的
spark.mllib不会突然消失 - 新项目、新功能建议优先基于
spark.ml来设计
六、性能优化:Breeze + 本地加速库的线性代数
机器学习离不开线性代数,而线性代数的效率直接决定了模型训练 / 推理的性能。
1. 线性代数依赖
MLlib 在底层使用了以下组件:
- Breeze:Scala 社区常用的数值计算 / 线性代数库
- dev.ludovic.netlib:负责在 JVM 与本地 BLAS 库之间做适配
这些组件可以在运行时调用系统中提供的本地加速库,例如:
- Intel MKL
- OpenBLAS 等
2. 为什么 Spark 不直接附带本地加速库?
官方说明很清楚:这些加速库属于系统级依赖,无法直接和 Spark 一起打包分发。
因此,MLlib 的策略是:
- 如果系统环境中配置了这些本地库,则启用加速路径
- 若未配置,则会退回到纯 JVM 实现
如果你在日志中看到类似:
WARNING: Failed to load implementation from:dev.ludovic.netlib.blas.JNIBLAS
就说明本地加速库没有加载成功,此时 MLlib 会使用纯 JVM 版本,性能会有所下降。
3. PySpark 端的依赖
在 Python 端使用 MLlib 时,还有一个小前提:
需要 NumPy 版本 1.4 及以上。
这主要是为了保证数值计算基础库的兼容性。
七、Spark 3.0 中 MLlib 的关键新特性
到 Spark 3.0 为止,MLlib 在 DataFrame API 方向已经非常成熟,同时加入了不少非常实用的增强能力。
下面挑重点拆解一下这些新特性。
1. 多列支持:特征变换一步到位
以下组件新增了 多列输入 支持:
BinarizerStringIndexerStopWordsRemover- PySpark 的
QuantileDiscretizer
这意味着:
- 你可以用一次 Transformer 调用处理多列特征
- 代码更简洁,Pipeline 更清晰
2. 基于树的特征变换(Tree-Based Feature Transformation)
新增了 Tree-Based Feature Transformation,可以通过决策树等模型:
- 自动学习特征之间的非线性交互关系
- 把这些信息编码为新的特征,供后续模型(如 LR、SVM)使用
这是一类在推荐、广告、风控中非常常见的套路——
树模型挖掘特征 + 线性模型做预测。
3. 新增 Evaluator:多标签和排序评估
增加了两个重要的评估器:
-
MultilabelClassificationEvaluator- 适用于多标签分类(一个样本可以有多个标签)
-
RankingEvaluator- 适用于推荐、搜索场景中的排序评估(如 NDCG、MAP 等)
对于推荐系统和信息检索场景,这是非常实用的补全能力。
4. 样本权重(Sample Weights)支持全面铺开
以下算法现在支持样本权重:
DecisionTreeClassifier/RegressorRandomForestClassifier/RegressorGBTClassifier/RegressorMulticlassClassificationEvaluatorRegressionEvaluatorBinaryClassificationEvaluatorBisectingKMeansKMeansGaussianMixture
样本权重在实际场景中的用途包括:
- 处理类别不平衡(给少数类更高权重)
- 对高价值样本(如转化用户)赋予更高权重
- 在抽样或去重场景中恢复原始分布
5. R 语言支持 PowerIterationClustering
为 R API 增加了 PowerIterationClustering 支持,
提升了 R 用户在 Spark 上做聚类分析的能力。
6. Spark ML Listener:跟踪 Pipeline 状态
新增的 Spark ML Listener 可以用来:
- 跟踪 ML Pipeline 的执行进度
- 记录各个阶段的状态与耗时
- 更好地做监控与排障
在生产环境中,这对于复杂流水线的可观测性非常重要。
7. GBT 支持验证集(Python)
在 Python API 中,Gradient Boosted Trees(GBT)新增了 验证集(validation set) 支持:
- 可以利用验证集做早停(early stopping)
- 降低过拟合风险
- 减少不必要的迭代,提高训练效率
8. RobustScaler:对异常值更鲁棒的特征缩放
RobustScaler 是新增的特征缩放器,特点是:
- 使用中位数和四分位距(IQR)做缩放
- 相比标准化,更不受极端异常值影响
适合金融、日志等大量噪声、离群点存在的场景。
9. Factorization Machines(分解机)分类与回归
Factorization Machines 在推荐系统、点击率预估中非常常见:
- 擅长建模稀疏特征之间的二阶交互
- 适合大量 One-Hot / ID 特征的广告、推荐场景
Spark 3.0 新增了:
- FM Classifier
- FM Regressor
为工业界常见的「FM + 深度模型」等组合提供了基础组件。
10. Gaussian / Complement Naive Bayes
朴素贝叶斯家族新增两个变种:
- Gaussian Naive Bayes:适合连续特征(假设特征服从高斯分布)
- Complement Naive Bayes:对类别不平衡更鲁棒,常用于文本多分类等场景
丰富了 Naive Bayes 在实际业务中的应用场景。
11. Scala / Python 功能对齐 + predictRaw / predictProbability 开放
Spark 3.0 还做了两方面的改善:
- Scala 和 Python 间的功能对齐:尽量避免「Scala 有、Python 没」的情况
- 在所有分类模型中将
predictRaw公共化,predictProbability也在绝大多数模型中对外公开(LinearSVCModel 除外)
这让你可以:
- 更方便地获取模型的原始打分 / 概率输出
- 自己做阈值调整、代价敏感决策、可解释性分析等工作
八、实战建议:现在做 Spark ML,应该怎么选型?
结合上面的内容,给出几个实践上的落地建议:
-
新项目 / 新系统:优先使用 DataFrame-based API(
spark.ml)- 与 SQL、Dataset、数据源高度统一
- 可以充分利用 Catalyst / Tungsten 优化
- 更适合构建易维护的 ML Pipeline
-
老系统:RDD API 继续用没问题,但可以规划渐进式迁移
- 不用一刀切重写
- 可以在新模块、新特性上优先使用 Pipeline + DataFrame
- 一边迭代业务,一边逐步收敛到 DataFrame 体系
-
性能层面:重视线性代数加速
- 尽量配置 Intel MKL / OpenBLAS 等本地库
- 注意集群各节点的一致性
- 留意
JNIBLAS相关日志,确认是否在走加速路径
-
善用 Spark 3.0+ 的高级特性
- 多列特征变换、样本权重、RobustScaler
- Factorization Machines、Naive Bayes 变体
- RankingEvaluator / MultilabelClassificationEvaluator 等评估器
九、总结
MLlib 从最初的一套 RDD 风格算法库,逐渐发展成了一个围绕 DataFrame、Pipeline 构建的 分布式机器学习平台组件。
- MLlib 没有被弃用,而是在不断演进
- RDD-based API 进入维护模式,但仍然可用
- DataFrame-based API(“Spark ML”)是 Spark 机器学习的现在和未来
如果你正在做基于 Spark 的机器学习平台建设或模型开发:
推荐的路径是:以
spark.ml为主,Pipeline 为核心,配合底层线性代数加速和 Spark 3.0+ 新特性,构建真正可落地、可扩展的分布式 ML 体系。
更多推荐
所有评论(0)