概述
MLlib 是 Spark 的机器学习(ML)库。其目标是使实用的机器学习可扩展且容易。在较高级别,它提供了以下工具:
- ML 算法:常见的学习算法,例如分类、回归、聚类和协作过滤
- 特征化:特征提取、变换、降维和选择
- 管道:用于构建、评估和调整 ML 管道的工具
- 持久性:保存和加载算法、模型和管道
- 实用程序:线性代数、统计信息、数据处理等
依赖库
MLlib 使用线性代数程序包 Breeze,该程序依赖于 netlib-java 进行优化的数值处理。
相关性
计算两个系列数据之间的相关性是统计中的常见操作。spark.ml 提供了很多系列中的灵活性,计算两两相关性。目前支持的相关方法是 Pearson 和 Spearman 的相关。
import org.apache.spark.ml.linalg.{Matrix, Vectors}
import org.apache.spark.ml.stat.Correlation
import org.apache.spark.sql.Row
val data = Seq(
Vectors.sparse(4, Seq((0, 1.0), (3, -2.0))),
Vectors.dense(4.0, 5.0, 0.0, 3.0),
Vectors.dense(6.0, 7.0, 0.0, 8.0),
Vectors.sparse(4, Seq((0, 9.0), (3, 1.0)))
)
val df = data.map(Tuple1.apply).toDF("features")
val Row(coeff1: Matrix) = Correlation.corr(df, "features").head
println(s"Pearson correlation matrix:\n $coeff1")
val Row(coeff2: Matrix) = Correlation.corr(df, "features", "spearman").head
println(s"Spearman correlation matrix:\n $coeff2")
补充:本地编程会报错
Seq没有toDF方法,应当在spark定义后,import spark.implicits._
假设检验
假设检验是一种强大的统计工具,可用来确定结果是否具有统计学意义,以及该结果是否偶然发生。spark.ml 目前支持 Pearson 的卡方(χ²)测试独立性。
ChiSquareTest 针对标签上的每个功能进行 Pearson 的独立性测试。对于每个特征,将(特征,标签)对转换为列联矩阵,针对该列矩阵计算卡方统计量。所有标签和特征值必须是分类的。
import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.stat.ChiSquareTest
val data = Seq(
(0.0, Vectors.dense(0.5, 10.0)),
(0.0, Vectors.dense(1.5, 20.0)),
(1.0, Vectors.dense(1.5, 30.0)),
(0.0, Vectors.dense(3.5, 30.0)),
(0.0, Vectors.dense(3.5, 40.0)),
(1.0, Vectors.dense(3.5, 40.0))
)
val df = data.toDF("label", "features")
val chi = ChiSquareTest.test(df, "features", "label").head
println(s"pValues = ${chi.getAs[Vector](0)}")
println(s"degreesOfFreedom ${chi.getSeq[Int](1).mkString("[", ",", "]")}")
println(s"statistics ${chi.getAs[Vector](2)}")
聚合运算
Summarizer 为 DataFrame 提供矢量列汇总统计。可用的度量是按列的最大值、最小值、平均值、方差和非零数,以及总数。
import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.stat.Summarizer
val data = Seq(
(Vectors.dense(2.0, 3.0, 5.0), 1.0),
(Vectors.dense(4.0, 6.0, 7.0), 2.0)
)
val df = data.toDF("features", "weight")
val (meanVal, varianceVal) = df.select(metrics("mean", "variance")
.summary($"features", $"weight").as("summary"))
.select("summary.mean", "summary.variance")
.as[(Vector, Vector)].first()
println(s"with weight: mean = ${meanVal}, variance = ${varianceVal}")
val (meanVal2, varianceVal2) = df.select(mean($"features"), variance($"features"))
.as[(Vector, Vector)].first()
println(s"without weight: mean = ${meanVal2}, sum = ${varianceVal2}")
管道中的主要概念
MLlib 对用于机器学习算法的 API 进行了标准化,从而使将多种算法组合到单个管道或工作流中变得更加容易。本节介绍了 Pipelines API 引入的关键概念,其中管道概念主要受 scikit-learn 项目的启发。
- DataFrame:此 ML API 使用 Spark SQL 的
DataFrame作为 ML 数据集,可以保存各种数据类型。例如,DataFrame可能有不同的列,用于存储文本、特征向量、真实标签和预测。 - Transformer:一种算法,可以将一个
DataFrame转换为另一个DataFrame。例如,ML 模型是一种Transformer,将具有特征的DataFrame转换为具有预测的DataFrame。 - Estimator:一种算法,可以拟合
DataFrame来产生Transformer。例如,学习算法是在DataFrame上训练并生成模型的Estimator。 - Pipeline:将多个
Transformer和Estimator链接在一起以指定 ML 工作流程。 - Parameter:所有
Transformer和Estimator共享一个用于指定参数的通用 API。
数据框
机器学习可以应用于多种数据类型,例如矢量、文本、图像和结构化数据。该 API 采用 Spark SQL 的 DataFrame,以支持多种数据类型。
DataFrame 支持许多基本和结构化类型;除了 Spark SQL 指南中列出的类型之外,DataFrame 还可以使用 ML Vector 类型。
DataFrame 可以从常规 RDD 直接或间接创建。中的列已命名,下面的代码示例使用诸如 "text"、"features" 和 "label" 之类的名称。
管道组件
Transformer
Transformer 是一个包含特征转换器和学习模型的抽象。从技术上讲,Transformer 实现了 transform() 方法,通常通过附加一个或多个列将 DataFrame 转换为另一种形式。例如:
- 特征转换器:可以采用
DataFrame,读取一列(例如文本),将其映射到新列(例如特征向量),并输出附加了映射列的新DataFrame - 学习模型:可以采用
DataFrame,读取包含特征向量的列,预测每个特征向量的标签,然后输出带有预测标签的新列
Estimator
Estimator 是学习算法或拟合数据的概念抽象。从技术上讲,Estimator 实现了一种 fit() 方法,该方法接受 DataFrame 并产生 Model,即 Transformer。例如,学习算法 LogisticRegression 是 Estimator,调用 fit() 训练出 LogisticRegressionModel,其为 Model,也是 Transformer。
管道组件的属性
Transformer.transform() 和 Estimator.fit() 都是无状态的。将来可能通过替代概念支持有状态算法。
Transformer 或 Estimator 的每个实例都有一个唯一的 ID,在指定参数时很有用。
Pipeline
在机器学习中,通常需要按顺序运行一系列算法来处理数据并从中学习。例如,一个简单的文本文档处理工作流程可能包括几个阶段:
- 将每个文档的文本拆分为单词
- 将每个文档的单词转换为数字特征向量
- 使用特征向量和标签学习预测模型
MLlib 将这样的工作流程表示为 Pipeline,其中包含要按特定顺序运行的一系列 PipelineStage(Transformer 和 Estimator)。
工作原理
Pipeline 被指定为一个阶段序列,每个阶段都是一个 Transformer 或一个 Estimator。这些阶段按顺序运行,输入 DataFrame 在通过每个阶段时都会进行转换:
- 对于
Transformer阶段,transform()方法作用于DataFrame - 对于
Estimator阶段,fit()作用于DataFrame,生成transform()方法
在训练时,Pipeline.fit() 在原始 DataFrame 上调用,经过各个阶段:
Tokenizer.transform()将原始文本文档拆分为单词,添加带有单词的新列HashingTF.transform()将 words 列转换为特征向量,添加带有向量新列的DataFrame- 由于
LogisticRegression是Estimator,Pipeline首先调用LogisticRegression.fit()产生LogisticRegressionModel,然后调用其transform()方法将DataFrame转换为另一个DataFrame
Pipeline 本身是一个 Estimator。运行 Pipeline 的 fit() 方法后,它会产生一个 PipelineModel,即一个 Transformer。PipelineModel 用于测试阶段,原始中的所有 Estimator 已变为 Transformer。
细节
- DAG Pipeline:
Pipeline被指定为一个有序数组。只要数据流图形成有向无环图(DAG),就可以创建非线性的Pipeline。当前基于每个阶段的输入和输出列名称隐式指定该图。 - 运行时检查:
Pipeline和PipelineModel在实际运行之前会进行运行时检查,使用DataFrame模式完成类型检查。 - 唯一的管道阶段:
Pipeline的阶段应该是唯一的实例。例如,同一个myHashingTF实例不应插入两次,但不同的实例myHashingTF1和myHashingTF2(都属于HashingTF类型)可以放入同一个Pipeline中。
参数
MLlib Estimator 和 Transformer 使用统一的 API 来指定参数。
Param是具有独立文件的命名参数ParamMap是一组(参数,值)对
将参数传递给算法的主要方式有两种:
- 设置实例的参数:例如
lr.setMaxIter(10),使lr.fit()最多使用 10 次迭代 - 传递
ParamMap给fit()或transform():ParamMap中的任何参数将覆盖先前通过 setter 方法指定的参数
参数属于 Estimator 和 Transformer 的特定实例。例如,如果两个 LogisticRegression 实例 lr1 和 lr2,则可以构建 ParamMap(lr1.maxIter -> 10, lr2.maxIter -> 20)。
代码示例
示例:Estimator、Transformer 和 Parameter
import org.apache.spark.ml.classification.LogisticRegression
import org.apache.spark.ml.linalg.{Vector, Vectors}
import org.apache.spark.ml.param.ParamMap
import org.apache.spark.sql.Row
// Prepare training data from a list of (label, features) tuples.
val training = spark.createDataFrame(Seq(
(1.0, Vectors.dense(0.0, 1.1, 0.1)),
(0.0, Vectors.dense(2.0, 1.0, -1.0)),
(0.0, Vectors.dense(2.0, 1.3, 1.0)),
(1.0, Vectors.dense(0.0, 1.2, -0.5))
)).toDF("label", "features")
// Create a LogisticRegression instance. This instance is an Estimator.
val lr = new LogisticRegression()
// Print out the parameters, documentation, and any default values.
println(s"LogisticRegression parameters:\n ${lr.explainParams()}\n")
// We may set parameters using setter methods.
lr.setMaxIter(10)
.setRegParam(0.01)
// Learn a LogisticRegression model. This uses the parameters stored in lr.
val model1 = lr.fit(training)
// Since model1 is a Model (i.e., a Transformer produced by an Estimator),
// we can view the parameters it used during fit().
// This prints the parameter (name: value) pairs, where names are unique IDs for this
// LogisticRegression instance.
println(s"Model 1 was fit using parameters: ${model1.parent.extractParamMap}")
// We may alternatively specify parameters using a ParamMap,
// which supports several methods for specifying parameters.
val paramMap = ParamMap(lr.maxIter -> 20)
.put(lr.maxIter, 30) // Specify 1 Param. This overwrites the original maxIter.
.put(lr.regParam -> 0.1, lr.threshold -> 0.55) // Specify multiple Params.
// One can also combine ParamMaps.
val paramMap2 = ParamMap(lr.probabilityCol -> "myProbability") // Change output column name.
val paramMapCombined = paramMap ++ paramMap2
// Now learn a new model using the paramMapCombined parameters.
// paramMapCombined overrides all parameters set earlier via lr.set* methods.
val model2 = lr.fit(training, paramMapCombined)
println(s"Model 2 was fit using parameters: ${model2.parent.extractParamMap}")
// Prepare test data.
val test = spark.createDataFrame(Seq(
(1.0, Vectors.dense(-1.0, 1.5, 1.3)),
(0.0, Vectors.dense(3.0, 2.0, -0.1)),
(1.0, Vectors.dense(0.0, 2.2, -1.5))
)).toDF("label", "features")
// Make predictions on test data using the Transformer.transform() method.
// LogisticRegression.transform will only use the 'features' column.
// Note that model2.transform() outputs a 'myProbability' column instead of the usual
// 'probability' column since we renamed the lr.probabilityCol parameter previously.
model2.transform(test)
.select("features", "label", "myProbability", "prediction")
.collect()
.foreach { case Row(features: Vector, label: Double, prob: Vector, prediction: Double) =>
println(s"($features, $label) -> prob=$prob, prediction=$prediction")
}
示例:Pipeline
import org.apache.spark.ml.{Pipeline, PipelineModel}
import org.apache.spark.ml.classification.LogisticRegression
import org.apache.spark.ml.feature.{HashingTF, Tokenizer}
import org.apache.spark.ml.linalg.Vector
import org.apache.spark.sql.Row
// Prepare training documents from a list of (id, text, label) tuples.
val training = spark.createDataFrame(Seq(
(0L, "a b c d e spark", 1.0),
(1L, "b d", 0.0),
(2L, "spark f g h", 1.0),
(3L, "hadoop mapreduce", 0.0)
)).toDF("id", "text", "label")
// Configure an ML pipeline, which consists of three stages: tokenizer, hashingTF, and lr.
val tokenizer = new Tokenizer()
.setInputCol("text")
.setOutputCol("words")
val hashingTF = new HashingTF()
.setNumFeatures(1000)
.setInputCol(tokenizer.getOutputCol)
.setOutputCol("features")
val lr = new LogisticRegression()
.setMaxIter(10)
.setRegParam(0.001)
val pipeline = new Pipeline()
.setStages(Array(tokenizer, hashingTF, lr))
// Fit the pipeline to training documents.
val model = pipeline.fit(training)
// Now we can optionally save the fitted pipeline to disk
model.write.overwrite().save("/tmp/spark-logistic-regression-model")
// We can also save this unfit pipeline to disk
pipeline.write.overwrite().save("/tmp/unfit-lr-model")
// And load it back in during production
val sameModel = PipelineModel.load("/tmp/spark-logistic-regression-model")
// Prepare test documents, which are unlabeled (id, text) tuples.
val test = spark.createDataFrame(Seq(
(4L, "spark i j k"),
(5L, "l m n"),
(6L, "spark hadoop spark"),
(7L, "apache hadoop")
)).toDF("id", "text")
// Make predictions on test documents.
model.transform(test)
.select("id", "text", "probability", "prediction")
.collect()
.foreach { case Row(id: Long, text: String, prob: Vector, prediction: Double) =>
println(s"($id, $text) --> prob=$prob, prediction=$prediction")
}
注意:如果使用
PipelineModel加载模型文件,可以通过stages(i)的方式获取管道中的某个环节的模型进行调用:val tokenizer = sameModel.stages(0)索引
i为模型文件目录的前缀。
特征提取与变换
TF-IDF
import org.apache.spark.ml.feature.{HashingTF, IDF, Tokenizer}
val sentenceData = spark.createDataFrame(Seq(
(0.0, "Hi I heard about Spark"),
(0.0, "I wish Java could use case classes"),
(1.0, "Logistic regression models are neat")
)).toDF("label", "sentence")
val tokenizer = new Tokenizer().setInputCol("sentence").setOutputCol("words")
val wordsData = tokenizer.transform(sentenceData)
val hashingTF = new HashingTF()
.setInputCol("words").setOutputCol("rawFeatures").setNumFeatures(20)
val featurizedData = hashingTF.transform(wordsData)
// alternatively, CountVectorizer can also be used to get term frequency vectors
val idf = new IDF().setInputCol("rawFeatures").setOutputCol("features")
val idfModel = idf.fit(featurizedData)
val rescaledData = idfModel.transform(featurizedData)
rescaledData.select("label", "features").show()
Word2Vec
import org.apache.spark.ml.feature.Word2Vec
import org.apache.spark.ml.linalg.Vector
import org.apache.spark.sql.Row
// Input data: Each row is a bag of words from a sentence or document.
val documentDF = spark.createDataFrame(Seq(
"Hi I heard about Spark".split(" "),
"I wish Java could use case classes".split(" "),
"Logistic regression models are neat".split(" ")
).map(Tuple1.apply)).toDF("text")
// Learn a mapping from words to Vectors.
val word2Vec = new Word2Vec()
.setInputCol("text")
.setOutputCol("result")
.setVectorSize(3)
.setMinCount(0)
val model = word2Vec.fit(documentDF)
val result = model.transform(documentDF)
result.collect().foreach { case Row(text: Seq[_], features: Vector) =>
println(s"Text: [${text.mkString(", ")}] => \nVector: $features\n") }
CountVectorizer
假设有以下 DataFrame,包含 id 和 texts 列:
| id | texts |
|---|---|
| 0 | Array(“a”, “b”, “c”) |
| 1 | Array(“a”, “b”, “b”, “c”, “a”) |
CountVectorizer.fit() 产生一个 CountVectorizerModel,词汇表为 (a, b, c)。转换后的输出列 vector 包含:
| id | texts | vector |
|---|---|---|
| 0 | Array(“a”, “b”, “c”) | (3,[0,1,2],[1.0,1.0,1.0]) |
| 1 | Array(“a”, “b”, “b”, “c”, “a”) | (3,[0,1,2],[2.0,2.0,1.0]) |
每个向量表示文档中 token 的计数。
import org.apache.spark.ml.feature.{CountVectorizer, CountVectorizerModel}
val df = spark.createDataFrame(Seq(
(0, Array("a", "b", "c")),
(1, Array("a", "b", "b", "c", "a"))
)).toDF("id", "words")
// fit a CountVectorizerModel from the corpus
val cvModel: CountVectorizerModel = new CountVectorizer()
.setInputCol("words")
.setOutputCol("features")
.setVocabSize(3)
.setMinDF(2)
.fit(df)
// alternatively, define CountVectorizerModel with a-priori vocabulary
val cvm = new CountVectorizerModel(Array("a", "b", "c"))
.setInputCol("words")
.setOutputCol("features")
cvModel.transform(df).show(false)
说明:在
feature包路径下,很多类是成对出现的,一般情况下,没有Model后缀的类用于模型训练,有Model后缀的类用于模型预测、加载、保存。例如:CountVectorizer训练后生成CountVectorizerModel,CountVectorizerModel可以load或save。带有Model后缀的类也可以设置setInputCol和setOutputCol。
FeatureHasher
FeatureHasher 将输入的多个特征通过 hash 函数进行向量化,返回结果为稀疏向量,稀疏向量中元素的值一般为整型数值,可以用于计算相似度。
假设有一个 DataFrame,包含 4 个输入列 real、bool、stringNum、string:
| real | bool | stringNum | string |
|---|---|---|---|
| 2.2 | true | 1 | foo |
| 3.3 | false | 2 | bar |
| 4.4 | false | 3 | baz |
| 5.5 | false | 4 | foo |
FeatureHasher.transform 的输出:
| real | bool | stringNum | string | features |
|---|---|---|---|---|
| 2.2 | true | 1 | foo | (262144,[51871, 63643,174475,253195],[1.0,1.0,2.2,1.0]) |
| 3.3 | false | 2 | bar | (262144,[6031, 80619,140467,174475],[1.0,1.0,1.0,3.3]) |
| 4.4 | false | 3 | baz | (262144,[24279,140467,174475,196810],[1.0,1.0,4.4,1.0]) |
| 5.5 | false | 4 | foo | (262144,[63643,140467,168512,174475],[1.0,1.0,1.0,5.5]) |
import org.apache.spark.ml.feature.FeatureHasher
val dataset = spark.createDataFrame(Seq(
(2.2, true, "1", "foo"),
(3.3, false, "2", "bar"),
(4.4, false, "3", "baz"),
(5.5, false, "4", "foo")
)).toDF("real", "bool", "stringNum", "string")
val hasher = new FeatureHasher()
.setInputCols("real", "bool", "stringNum", "string")
.setOutputCol("features")
val featurized = hasher.transform(dataset)
featurized.show(false)
可以使用
setNumFeatures(4)控制返回的稀疏向量长度,稀疏向量长度会限制可以表示的值空间大小。
特征变换器
Tokenizer
import org.apache.spark.ml.feature.{RegexTokenizer, Tokenizer}
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
val sentenceDataFrame = spark.createDataFrame(Seq(
(0, "Hi I heard about Spark"),
(1, "I wish Java could use case classes"),
(2, "Logistic,regression,models,are,neat")
)).toDF("id", "sentence")
val tokenizer = new Tokenizer().setInputCol("sentence").setOutputCol("words")
val regexTokenizer = new RegexTokenizer()
.setInputCol("sentence")
.setOutputCol("words")
.setPattern("\\W") // alternatively .setPattern("\\w+").setGaps(false)
val countTokens = udf { (words: Seq[String]) => words.length }
val tokenized = tokenizer.transform(sentenceDataFrame)
tokenized.select("sentence", "words")
.withColumn("tokens", countTokens(col("words"))).show(false)
val regexTokenized = regexTokenizer.transform(sentenceDataFrame)
regexTokenized.select("sentence", "words")
.withColumn("tokens", countTokens(col("words"))).show(false)
可以使用
setStopWords(Array("word1", "word2", ...))设置停用词。
StopWordsRemover
假设有以下 DataFrame:
| id | raw |
|---|---|
| 0 | [I, saw, the, red, baloon] |
| 1 | [Mary, had, a, little, lamb] |
应用 StopWordsRemover 后:
| id | raw | filtered |
|---|---|---|
| 0 | [I, saw, the, red, baloon] | [saw, red, baloon] |
| 1 | [Mary, had, a, little, lamb] | [Mary, little, lamb] |
停用词 "I"、"the"、"had"、"a" 已被过滤。
val remover = new StopWordsRemover()
.setInputCol("raw")
.setOutputCol("filtered")
val dataSet = spark.createDataFrame(Seq(
(0, Seq("I", "saw", "the", "red", "balloon")),
(1, Seq("Mary", "had", "a", "little", "lamb"))
)).toDF("id", "raw")
remover.transform(dataSet).show(false)
N-gram
NGram 用于生成 n-gram 组合,例如,指定 N=2 时,Array("Hi", "I", "heard", "about", "Spark") 会返回 Array("Hi I", "I heard", "heard about", "about Spark")。
import org.apache.spark.ml.feature.NGram
val wordDataFrame = spark.createDataFrame(Seq(
(0, Array("Hi", "I", "heard", "about", "Spark")),
(1, Array("I", "wish", "Java", "could", "use", "case", "classes")),
(2, Array("Logistic", "regression", "models", "are", "neat"))
)).toDF("id", "words")
val ngram = new NGram().setN(2).setInputCol("words").setOutputCol("ngrams")
val ngramDataFrame = ngram.transform(wordDataFrame)
ngramDataFrame.select("ngrams").show(false)
Binarizer
对数值类型进行二值化处理,通过指定阈值,大于阈值部分为 1,小于等于阈值部分为 0。
val data = Array((0, 0.1), (1, 0.8), (2, 0.2))
val dataFrame = spark.createDataFrame(data).toDF("id", "feature")
val binarizer: Binarizer = new Binarizer()
.setInputCol("feature")
.setOutputCol("binarized_feature")
.setThreshold(0.5)
val binarizedDataFrame = binarizer.transform(dataFrame)
println(s"Binarizer output with Threshold = ${binarizer.getThreshold}")
binarizedDataFrame.show()
通过
setThreshold(0.5)将阈值指定为0.5,当数值等于0.5时为0。
PCA
PCA 降维,可以作为一种特征选择器来使用。
import org.apache.spark.ml.feature.PCA
import org.apache.spark.ml.linalg.Vectors
val data = Array(
Vectors.sparse(5, Seq((1, 1.0), (3, 7.0))),
Vectors.dense(2.0, 0.0, 3.0, 4.0, 5.0),
Vectors.dense(4.0, 0.0, 0.0, 6.0, 7.0)
)
val df = spark.createDataFrame(data.map(Tuple1.apply)).toDF("features")
val pca = new PCA()
.setInputCol("features")
.setOutputCol("pcaFeatures")
.setK(3)
.fit(df)
val result = pca.transform(df).select("pcaFeatures")
result.show(false)
PolynomialExpansion(多项式扩展)
对数值型特征进行多项式扩展,将原始特征组合转化为多项式形式,从而增加特征复杂性,有助于捕获更多非线性关系。
import org.apache.spark.ml.feature.PolynomialExpansion
import org.apache.spark.ml.linalg.Vectors
val data = Array(
Vectors.dense(2.0, 1.0),
Vectors.dense(0.0, 0.0),
Vectors.dense(3.0, -1.0)
)
val df = spark.createDataFrame(data.map(Tuple1.apply)).toDF("features")
val polyExpansion = new PolynomialExpansion()
.setInputCol("features")
.setOutputCol("polyFeatures")
.setDegree(3)
val polyDF = polyExpansion.transform(df)
polyDF.show(false)
setDegree(3)表示生成原始特征的 3 次方项以及原始特征之间的交叉项,多项式形式为:x³ + x²y + xy² + y³ + x² + xy + y² + x + ysetDegree(1)表示生成的多项式形式为:x + y,例如(2.0, 1.0)对应的多项式为2·x + 1·y
Discrete Cosine Transform(DCT,离散余弦变换)
import org.apache.spark.ml.linalg.Vectors
val data = Seq(
Vectors.dense(0.0, 1.0, -2.0, 3.0),
Vectors.dense(-1.0, 2.0, 4.0, -7.0),
Vectors.dense(14.0, -2.0, -5.0, 1.0))
val df = spark.createDataFrame(data.map(Tuple1.apply)).toDF("features")
val dct = new DCT()
.setInputCol("features")
.setOutputCol("featuresDCT")
.setInverse(false)
val dctDf = dct.transform(df)
dctDf.select("featuresDCT").show(false)
StringIndexer
将字符串列转换为数值索引列。
假设有以下 DataFrame:
| id | category |
|---|---|
| 0 | a |
| 1 | b |
| 2 | c |
| 3 | a |
| 4 | a |
| 5 | c |
应用 StringIndexer 后:
| id | category | categoryIndex |
|---|---|---|
| 0 | a | 0.0 |
| 1 | b | 2.0 |
| 2 | c | 1.0 |
| 3 | a | 0.0 |
| 4 | a | 0.0 |
| 5 | c | 1.0 |
"a" 得到索引 0,因为它是最频繁的;其次是索引 1 的 "c";索引 2 的 "b"。
处理未知标签的策略
对于在一个数据集上拟合 StringIndexer,然后使用它转换另一个数据集时,有三种方式处理未知标签:
- error(默认):抛出异常
- skip:跳过包含未知标签的行
- keep:将未知标签放入一个特殊的额外桶中,索引为
numLabels
import org.apache.spark.ml.feature.StringIndexer
val df = spark.createDataFrame(
Seq((0, "a"), (1, "b"), (2, "c"), (3, "a"), (4, "a"), (5, "c"))
).toDF("id", "category")
val indexer = new StringIndexer()
.setInputCol("category")
.setOutputCol("categoryIndex")
val indexed = indexer.fit(df).transform(df)
indexed.show()
setHandleInvalid(value: String):设置输入数据中的无效值(未见过的分类)的处理方式setStringOrderType(value: String):设置字符串的排序方式,frequencyDesc按频率降序排列,frequencyAsc按频率升序排列
IndexToString
在 StringIndexer 的基础上,将数值索引还原为原始字符串标签。
假设有以下 DataFrame:
| id | categoryIndex |
|---|---|
| 0 | 0.0 |
| 1 | 2.0 |
| 2 | 1.0 |
| 3 | 0.0 |
| 4 | 0.0 |
| 5 | 1.0 |
应用 IndexToString 后:
| id | categoryIndex | originalCategory |
|---|---|---|
| 0 | 0.0 | a |
| 1 | 2.0 | b |
| 2 | 1.0 | c |
| 3 | 0.0 | a |
| 4 | 0.0 | a |
| 5 | 1.0 | c |
import org.apache.spark.ml.attribute.Attribute
import org.apache.spark.ml.feature.{IndexToString, StringIndexer}
val df = spark.createDataFrame(Seq(
(0, "a"),
(1, "b"),
(2, "c"),
(3, "a"),
(4, "a"),
(5, "c")
)).toDF("id", "category")
val indexer = new StringIndexer()
.setInputCol("category")
.setOutputCol("categoryIndex")
.fit(df)
val indexed = indexer.transform(df)
println(s"Transformed string column '${indexer.getInputCol}' " +
s"to indexed column '${indexer.getOutputCol}'")
indexed.show()
val inputColSchema = indexed.schema(indexer.getOutputCol)
println(s"StringIndexer will store labels in output column metadata: " +
s"${Attribute.fromStructField(inputColSchema).toString}\n")
val converter = new IndexToString()
.setInputCol("categoryIndex")
.setOutputCol("originalCategory")
val converted = converter.transform(indexed)
println(s"Transformed indexed column '${converter.getInputCol}' back to original string " +
s"column '${converter.getOutputCol}' using labels in metadata")
converted.select("id", "categoryIndex", "originalCategory").show()
df.schema("colName")用于返回列的StructField信息Attribute.fromStructField用于查看StringIndexer存储在输出列元数据中的标签信息df.schema("colName").metadata也可以达到相同的效果setLabel(value: Array[String]):设置新的类型名称,索引地址从小到大依次转换为指定的类型名称converter实际上使用了indexed.schema("categoryIndex").metadata中保存的元数据信息,用于恢复categoryIndex到originalCategory的映射关系
OneHotEncoder
独热编码仅支持数值类型,将单列数值类型转换为稀疏向量。如果需要对多列组合进行独热编码,一种简单的方式是对多列分别进行独热编码,然后使用 VectorAssembler 合并多个独热编码的列。
import org.apache.spark.ml.feature.OneHotEncoderEstimator
val df = spark.createDataFrame(Seq(
(0.0, 1.0),
(1.0, 0.0),
(2.0, 1.0),
(0.0, 2.0),
(0.0, 1.0),
(2.0, 0.0)
)).toDF("categoryIndex1", "categoryIndex2")
val encoder = new OneHotEncoderEstimator()
.setInputCols(Array("categoryIndex1", "categoryIndex2"))
.setOutputCols(Array("categoryVec1", "categoryVec2"))
val model = encoder.fit(df)
val encoded = model.transform(df)
encoded.show()
OneHotEncoderEstimator指定多个输入列时会产生多个输出列,如果需要将多个输入列合并成一种特征,则需要提前进行VectorAssembler处理setDropLast设置是否丢弃最后一个类别,默认为true。如果设置为false,则生成每个类别的 OneHot 编码。例如对于特征类别1、2、3,对应 OneHot 编码为[1,0]、[0,1]、[0,0],2 位编码已足以表达 3 个类别categorySizes获取每个分类特征的类别数量(拟合模型之后才可用)
VectorIndexer
将特征向量转换为整型值索引,方便检索。
libsvm 格式说明: libsvm 是常用的存储稀疏矩阵的文件类型,格式如下:
<label> <index1>:<value1> <index2>:<value2> ...
label:数据的标签或类别index1:特征的索引,从 1 开始,必须按顺序排列。注意:解析为稀疏向量后,稀疏向量索引从 0 开始value1:特征的值,仅记录非零值libsvm 文件样例:
1 1:2.1 2:3.5 4:1.2 0 1:1.1 3:2.0 4:3.2 1 2:4.2 3:1.8 5:2.5
import org.apache.spark.ml.feature.VectorIndexer
val data = spark.read.format("libsvm").load("libsvm.txt")
val indexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexed")
.setMaxCategories(10)
val indexerModel = indexer.fit(data)
val categoricalFeatures: Set[Int] = indexerModel.categoryMaps.keys.toSet
println(s"Chose ${categoricalFeatures.size} " +
s"categorical features: ${categoricalFeatures.mkString(", ")}")
// Create new column "indexed" with categorical values transformed to indices
val indexedData = indexerModel.transform(data)
indexedData.show()
注意事项:
setMaxCategories(10)是指元数据中能够容纳的最大索引数量为 10 个,即 0-9- 实际使用的最大索引数量和
features列中的 Vector 长度一致,矩阵 Vector 长度为 5,因此最大索引数量为 5,即 0-4- 如果设置的最大索引长度小于实际需要解析的 Vector 数量,则超出的部分会被直接丢弃。如
setMaxCategories(2),则只会在元数据信息中保存 0 和 1 两个索引,使用indexerModel进行transform时只转换 0 和 1 对应的值,其他值按原值输出- 需要将 Index 还原为 Vector 时,可以通过
indexerModel.categoryMaps提供的映射关系进行还原indexedData中的indexed列带有元数据信息,包含了索引值与数值的对应关系
Interaction
对多个向量取笛卡尔积,计算元素两两乘积。
假设输入列 id1、vec1、vec2,Interaction 输出列 interactedCol:
| id1 | vec1 | vec2 | interactedCol |
|---|---|---|---|
| 1 | [1.0,2.0,3.0] | [8.0,4.0,5.0] | [8.0,4.0,5.0,16.0,8.0,10.0,24.0,12.0,15.0] |
import org.apache.spark.ml.feature.Interaction
import org.apache.spark.ml.feature.VectorAssembler
val df = spark.createDataFrame(Seq(
(1, 1, 2, 3, 8, 4, 5),
(2, 4, 3, 8, 7, 9, 8),
(3, 6, 1, 9, 2, 3, 6),
(4, 10, 8, 6, 9, 4, 5),
(5, 9, 2, 7, 10, 7, 3),
(6, 1, 1, 4, 2, 8, 4)
)).toDF("id1", "id2", "id3", "id4", "id5", "id6", "id7")
// VectorAssembler 相当于对多个向量取并集
val assembler1 = new VectorAssembler().
setInputCols(Array("id2", "id3", "id4")).
setOutputCol("vec1")
val assembled1 = assembler1.transform(df)
val assembler2 = new VectorAssembler().
setInputCols(Array("id5", "id6", "id7")).
setOutputCol("vec2")
val assembled2 = assembler2.transform(assembled1).select("id1", "vec1", "vec2")
// Interaction 对多个向量取笛卡尔积,获取元素两两乘积
val interaction = new Interaction()
.setInputCols(Array("id1", "vec1", "vec2"))
.setOutputCol("interactedCol")
val interacted = interaction.transform(assembled2)
interacted.show(truncate = false)
计算规则:样例数据
id1[d]、vec1[a1, b1, c1]、vec2[a2, b2, c2],计算结果:[d*a1*a2, d*a1*b2, d*a1*c2, d*b1*a2, d*b1*b2, d*b1*c2, d*c1*a2, d*c1*b2, d*c1*c2]
VectorAssembler:相当于对多个向量取并集,将多个输入列合并为一个向量列。
数据标准化与归一化
Normalizer(归一化)
将特征向量的每个分量缩放到单位范数(即 Lᵖ 范数,其中 p 是可配置参数)。
可以不进行模型训练。
import org.apache.spark.ml.feature.Normalizer
import org.apache.spark.ml.linalg.Vectors
val dataFrame = spark.createDataFrame(Seq(
(0, Vectors.dense(1.0, 0.5, -1.0)),
(1, Vectors.dense(2.0, 1.0, 1.0)),
(2, Vectors.dense(4.0, 10.0, 2.0))
)).toDF("id", "features")
// Normalize each Vector using L¹ norm.
val normalizer = new Normalizer()
.setInputCol("features")
.setOutputCol("normFeatures")
.setP(1.0)
val l1NormData = normalizer.transform(dataFrame)
println("Normalized using L¹ norm")
l1NormData.show()
// Normalize each Vector using L^∞ norm.
val lInfNormData = normalizer.transform(dataFrame, normalizer.p -> Double.PositiveInfinity)
println("Normalized using L^inf norm")
lInfNormData.show()
- L¹ 范数归一化:每个分量除以向量的 L¹ 范数(绝对值之和),即
||x||₁ = |x₁| + |x₂| + ... + |xₙ|- L² 范数归一化:每个分量除以向量的 L² 范数(平方和),即
||x||₂ = √(x₁² + x₂² + ... + xₙ²)setP(1.0)表示每个元素除以 L¹ 范数,setP(2.0)表示每个元素除以 L² 范数
StandardScaler(标准化)
每个特征值减去均值,除以标准差,使特征向量的均值为 0、标准差为 1。
可以不进行模型训练。
import org.apache.spark.ml.feature.StandardScaler
val dataFrame = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
val scaler = new StandardScaler()
.setInputCol("features")
.setOutputCol("scaledFeatures")
.setWithStd(true)
.setWithMean(false)
// Compute summary statistics by fitting the StandardScaler.
val scalerModel = scaler.fit(dataFrame)
// Normalize each feature to have unit standard deviation.
val scaledData = scalerModel.transform(dataFrame)
scaledData.show()
setWithStd(true):特征值将按照标准差进行缩放,确保方差等于 1setWithMean(true):特征值将按照均值进行中心化,使特征向量均值为 0
MinMaxScaler(最大最小标准化)
减去最小值并除以最大值减最小值的差,缩放到指定范围。
可以不进行模型训练。
import org.apache.spark.ml.feature.MinMaxScaler
import org.apache.spark.ml.linalg.Vectors
val dataFrame = spark.createDataFrame(Seq(
(0, Vectors.dense(1.0, 0.1, -1.0)),
(1, Vectors.dense(2.0, 1.1, 1.0)),
(2, Vectors.dense(3.0, 10.1, 3.0))
)).toDF("id", "features")
val scaler = new MinMaxScaler()
.setInputCol("features")
.setOutputCol("scaledFeatures")
// Compute summary statistics and generate MinMaxScalerModel
val scalerModel = scaler.fit(dataFrame)
// rescale each feature to range [min, max].
val scaledData = scalerModel.transform(dataFrame)
println(s"Features scaled to range: [${scaler.getMin}, ${scaler.getMax}]")
scaledData.select("features", "scaledFeatures").show()
setMin设置标准化后的最小值,setMax设置标准化后的最大值。通常setMin(0.0)和setMax(1.0),表示将特征值缩放到[0,1]范围内。
MaxAbsScaler(最大绝对值缩放)
将特征向量进行缩放,确保每个特征的绝对值在最大绝对值内,默认将特征缩放到 [-1, 1],适用于保留特征值的符号和相对比例。
可以不进行模型训练。
import org.apache.spark.ml.feature.MaxAbsScaler
import org.apache.spark.ml.linalg.Vectors
val dataFrame = spark.createDataFrame(Seq(
(0, Vectors.dense(1.0, 0.1, -8.0)),
(1, Vectors.dense(2.0, 1.0, -4.0)),
(2, Vectors.dense(4.0, 10.0, 8.0))
)).toDF("id", "features")
val scaler = new MaxAbsScaler()
.setInputCol("features")
.setOutputCol("scaledFeatures")
// Compute summary statistics and generate MaxAbsScalerModel
val scalerModel = scaler.fit(dataFrame)
// rescale each feature to range [-1, 1]
val scaledData = scalerModel.transform(dataFrame)
scaledData.select("features", "scaledFeatures").show()
setMaxAbs(1.0)设置最大绝对值,将特征值放缩至[-1, 1]。
Bucketizer(数据分箱)
连续数值离散化,将连续值映射到离散的桶中。
import org.apache.spark.ml.feature.Bucketizer
val splits = Array(Double.NegativeInfinity, -0.5, 0.0, 0.5, Double.PositiveInfinity)
val data = Array(-999.9, -0.5, -0.3, 0.0, 0.2, 999.9)
val dataFrame = spark.createDataFrame(data.map(Tuple1.apply)).toDF("features")
val bucketizer = new Bucketizer()
.setInputCol("features")
.setOutputCol("bucketedFeatures")
.setSplits(splits)
// Transform original data into its bucket index.
val bucketedData = bucketizer.transform(dataFrame)
println(s"Bucketizer output with ${bucketizer.getSplits.length-1} buckets")
bucketedData.show()
// 多列分箱
val splitsArray = Array(
Array(Double.NegativeInfinity, -0.5, 0.0, 0.5, Double.PositiveInfinity),
Array(Double.NegativeInfinity, -0.3, 0.0, 0.3, Double.PositiveInfinity))
val data2 = Array(
(-999.9, -999.9),
(-0.5, -0.2),
(-0.3, -0.1),
(0.0, 0.0),
(0.2, 0.4),
(999.9, 999.9))
val dataFrame2 = spark.createDataFrame(data2).toDF("features1", "features2")
val bucketizer2 = new Bucketizer()
.setInputCols(Array("features1", "features2"))
.setOutputCols(Array("bucketedFeatures1", "bucketedFeatures2"))
.setSplitsArray(splitsArray)
// Transform original data into its bucket index.
val bucketedData2 = bucketizer2.transform(dataFrame2)
println(s"Bucketizer output with [" +
s"${bucketizer2.getSplitsArray(0).length-1}, " +
s"${bucketizer2.getSplitsArray(1).length-1}] buckets for each input column")
bucketedData2.show()
ElementwiseProduct(元素级乘法)
执行元素级乘法操作,将每个特征向量的每个元素与一个相应的权重进行乘法操作,用于数据预处理过程中对特征进行加权或缩放。
可以不进行模型训练。
import org.apache.spark.ml.feature.ElementwiseProduct
import org.apache.spark.ml.linalg.Vectors
// Create some vector data; also works for sparse vectors
val dataFrame = spark.createDataFrame(Seq(
("a", Vectors.dense(1.0, 2.0, 3.0)),
("b", Vectors.dense(4.0, 5.0, 6.0)))).toDF("id", "vector")
val transformingVector = Vectors.dense(0.0, 1.0, 2.0)
val transformer = new ElementwiseProduct()
.setScalingVec(transformingVector)
.setInputCol("vector")
.setOutputCol("transformedVector")
// Batch transform the vectors to create new column:
transformer.transform(dataFrame).show()
setScalingVec(transformingVector)设置权重向量,让vector列中的元素与transformingVector向量中对应元素相乘。
SQLTransformer
SQL 转换器,用于执行 SQL 语句,输入的 DataFrame 在 SQL 语句中用 __THIS__ 代替。
假设有以下 DataFrame:
| id | v1 | v2 |
|---|---|---|
| 0 | 1.0 | 3.0 |
| 2 | 2.0 | 5.0 |
使用 SQLTransformer 执行语句 "SELECT *, (v1 + v2) AS v3, (v1 * v2) AS v4 FROM __THIS__" 后:
| id | v1 | v2 | v3 | v4 |
|---|---|---|---|---|
| 0 | 1.0 | 3.0 | 4.0 | 3.0 |
| 2 | 2.0 | 5.0 | 7.0 | 10.0 |
import org.apache.spark.ml.feature.SQLTransformer
val df = spark.createDataFrame(
Seq((0, 1.0, 3.0), (2, 2.0, 5.0))).toDF("id", "v1", "v2")
val sqlTrans = new SQLTransformer().setStatement(
"SELECT *, (v1 + v2) AS v3, (v1 * v2) AS v4 FROM __THIS__")
sqlTrans.transform(df).show()
VectorAssembler
VectorAssembler 相当于对多个向量取并集,将多个输入列合并为一个向量列。
假设有以下 DataFrame:
| id | hour | mobile | userFeatures | clicked |
|---|---|---|---|---|
| 0 | 18 | 1.0 | [0.0, 10.0, 0.5] | 1.0 |
将 hour、mobile、userFeatures 组合成 features 列后:
| id | hour | mobile | userFeatures | clicked | features |
|---|---|---|---|---|---|
| 0 | 18 | 1.0 | [0.0, 10.0, 0.5] | 1.0 | [18.0, 1.0, 0.0, 10.0, 0.5] |
import org.apache.spark.ml.feature.VectorAssembler
import org.apache.spark.ml.linalg.Vectors
val dataset = spark.createDataFrame(
Seq((0, 18, 1.0, Vectors.dense(0.0, 10.0, 0.5), 1.0))
).toDF("id", "hour", "mobile", "userFeatures", "clicked")
val assembler = new VectorAssembler()
.setInputCols(Array("hour", "mobile", "userFeatures"))
.setOutputCol("features")
val output = assembler.transform(dataset)
println("Assembled columns 'hour', 'mobile', 'userFeatures' to vector column 'features'")
output.select("features", "clicked").show(false)
自定义 Vector 类型列
注意:PySpark 和 Spark 有一定区别,PySpark 使用
pyspark.ml.linalg.VectorUDT类型,Spark 使用org.apache.spark.ml.linalg.SQLDataTypes.VectorType。注意:数据初始化使用
Vectors.dense,PySpark 初始化的数据使用元组封装,Spark 初始化的数据使用Row封装,否则最终结果为null。注意:PySpark 中,如果在 UDF 中对
Vector进行操作,且不指定返回类型为VectorUDT(),一定会报错:net.razorvine.pickle.PickleException: expected zero arguments for construction of ClassDict (for pyspark.ml.linalg.DenseVector),主要原因为 Python 对一些对象的序列化支持不够好。
PySpark 版本:
from pyspark.ml.linalg import VectorUDT, Vectors
from pyspark.sql.types import *
schema = StructType([StructField('vector', VectorUDT())])
df = spark.createDataFrame([(Vectors.dense(1.0, 2.0, 3.0), )], schema)
Spark 版本:
import org.apache.spark.ml.linalg.SQLDataTypes.VectorType
import org.apache.spark.ml.linalg.Vectors
import org.apache.spark.sql.types._
val schema = StructType(StructField("vec", VectorType) :: Nil)
val df = spark.createDataFrame(sc.parallelize(List(Row(Vectors.dense(1, 2, 3.0)))), schema)
Vector 类型列的通用解析方法
注意:Vector 支持索引,但是索引后的数值类型为 numpy 类型,无法序列化,需要强制转型为 Python 内置类型。
注意:UDF 需要指定返回类型,返回类型必须与 Vector 元素强转之后的类型相对应,否则结果为
null。
PySpark 版本:
vectorToInt = udf(lambda vector: int(vector[0]), IntegerType())
vectorToStr = udf(lambda vector: str(vector[1]), StringType())
vectorToFloat = udf(lambda vector: float(vector[2]), FloatType())
df.withColumn("int", vectorToInt(col("vector")))\
.withColumn("str", vectorToStr(col("vector")))\
.withColumn("float", vectorToFloat(col("vector"))).show(truncate=False)
Spark 版本:
import org.apache.spark.ml.linalg.Vector
val vectorToInt = udf((vector: Vector) => vector(0).asInstanceOf[Int], IntegerType)
val vectorToStr = udf((vector: Vector) => String.valueOf(vector(1)), StringType)
val vectorToFloat = udf((vector: Vector) => vector(2).asInstanceOf[Float], FloatType)
df.withColumn("int", vectorToInt(col("vec")))
.withColumn("str", vectorToStr(col("vec")))
.withColumn("float", vectorToFloat(col("vec"))).show(false)
VectorSizeHint
为特征向量的大小提供提示或约束,用于确保特征向量的大小与模型的期望输入大小一致,常用于处理具有可变维度特征的数据集,可以对特征向量大小不符合要求的行进行过滤,确保后续流程正常进行。
不需要模型训练。
有时显式地为 VectorType 列指定向量的大小是很有用的。例如,VectorAssembler 使用其输入列中的大小信息为其输出列生成大小信息和元数据。虽然在某些情况下,可以通过检查列的内容来获取此信息,但在流式 DataFrame 中,直到流启动后,内容才可用。VectorSizeHint 允许用户显式指定列的向量大小。
要使用 VectorSizeHint,用户必须设置 inputCol 和 size 参数。将此转换器应用于 DataFrame 将生成一个新的 DataFrame,其中包含指定向量大小的 inputCol 的更新元数据。
VectorSizeHint 还可以采用可选的 handleInvalid 参数,控制向量列包含 null 或大小错误的向量时的行为:
- error(默认):引发异常
- skip:包含无效值的行应该从结果 DataFrame 中筛选出来
- optimistic:不检查列中是否存在无效值,保留所有行(注意:可能会导致不一致状态)
import org.apache.spark.ml.feature.{VectorAssembler, VectorSizeHint}
import org.apache.spark.ml.linalg.Vectors
val dataset = spark.createDataFrame(
Seq(
(0, 18, 1.0, Vectors.dense(0.0, 10.0, 0.5), 1.0),
(0, 18, 1.0, Vectors.dense(0.0, 10.0), 0.0))
).toDF("id", "hour", "mobile", "userFeatures", "clicked")
val sizeHint = new VectorSizeHint()
.setInputCol("userFeatures")
.setHandleInvalid("skip")
.setSize(3)
val datasetWithSize = sizeHint.transform(dataset)
println("Rows where 'userFeatures' is not the right size are filtered out")
datasetWithSize.show(false)
val assembler = new VectorAssembler()
.setInputCols(Array("hour", "mobile", "userFeatures"))
.setOutputCol("features")
// This dataframe can be used by downstream transformers as before
val output = assembler.transform(datasetWithSize)
println("Assembled columns 'hour', 'mobile', 'userFeatures' to vector column 'features'")
output.select("features", "clicked").show(false)
QuantileDiscretizer
QuantileDiscretizer 获取具有连续特征的列,并输出具有二进制分类特征的列。箱子的数量由 numBuckets 参数设置。如果输入的不同值太少,无法创建足够的不同分位数,则使用的桶数可能小于此值。
NaN 值处理:在 QuantileDiscretizer 拟合期间,将从列中删除 NaN 值。在转换过程中,Bucketizer 在数据集中找到 NaN 值时会引发错误,但用户可以通过设置 handleInvalid 来选择保留或删除 NaN 值。如果选择保留 NaN 值,会将其放入特殊的桶中(例如,4 个桶时,NaN 放入桶 [4])。
算法:使用近似算法选择 bin 范围(关于 approxQuantile 的详细描述,请参阅文档)。近似的精度可以用相对误差参数来控制。当设置为零时,计算精确分位数(计算精确分位数是一项昂贵的操作)。
假设有以下 DataFrame:
| id | hour |
|---|---|
| 0 | 18.0 |
| 1 | 19.0 |
| 2 | 8.0 |
| 3 | 5.0 |
| 4 | 2.2 |
设置 numBuckets = 3 后:
| id | hour | result |
|---|---|---|
| 0 | 18.0 | 2.0 |
| 1 | 19.0 | 2.0 |
| 2 | 8.0 | 1.0 |
| 3 | 5.0 | 1.0 |
| 4 | 2.2 | 0.0 |
import org.apache.spark.ml.feature.QuantileDiscretizer
val data = Array((0, 18.0), (1, 19.0), (2, 8.0), (3, 5.0), (4, 2.2))
val df = spark.createDataFrame(data).toDF("id", "hour")
val discretizer = new QuantileDiscretizer()
.setInputCol("hour")
.setOutputCol("result")
.setNumBuckets(3)
val result = discretizer.fit(df).transform(df)
result.show(false)
Imputer
用于填充数据中的缺失值,可以指定使用均值或中位数进行填充。
输入列应为 DoubleType 或 FloatType。当前 Imputer 不支持分类特征,并且可能为包含分类特征的列创建不正确的值。Imputer 可以通过 .setMissingValue 填充除 NaN 以外的自定义值。例如 .setMissingValue(0) 将填充所有出现的 0。
注:输入列中的所有
null值都被视为缺失,因此也进行了插补。
假设有以下 DataFrame:
| a | b |
|---|---|
| 1.0 | Double.NaN |
| 2.0 | Double.NaN |
| Double.NaN | 3.0 |
| 4.0 | 4.0 |
| 5.0 | 5.0 |
列 a 和列 b 的代理值分别为 3.0 和 4.0。转换后:
| a | b | out_a | out_b |
|---|---|---|---|
| 1.0 | Double.NaN | 1.0 | 4.0 |
| 2.0 | Double.NaN | 2.0 | 4.0 |
| Double.NaN | 3.0 | 3.0 | 3.0 |
| 4.0 | 4.0 | 4.0 | 4.0 |
| 5.0 | 5.0 | 5.0 | 5.0 |
import org.apache.spark.ml.feature.Imputer
val df = spark.createDataFrame(Seq(
(1.0, Double.NaN),
(2.0, Double.NaN),
(Double.NaN, 3.0),
(4.0, 4.0),
(5.0, 5.0)
)).toDF("a", "b")
val imputer = new Imputer()
.setInputCols(Array("a", "b"))
.setOutputCols(Array("out_a", "out_b"))
val model = imputer.fit(df)
model.transform(df).show()
setMissingValue(0.0)指定用0替换缺失值setStrategy("mean")指定用均值替换缺失值,可选值:mean(均值)、median(中位数)、mode(众数)
特征选择器
VectorSlicer
矢量切片器是一个转换器,它接受一个特征向量并输出一个新的特征向量和一个子数组的原始特征。它对于从向量列中提取特征非常有用。
VectorSlicer 接受具有指定索引的向量列,然后输出一个新的向量列,有两种类型的索引:
- 整数索引:
setIndices(),表示向量中的索引 - 字符串索引:
setNames(),表示向量中的功能名称(要求 vector 列具有AttributeGroup)
可以同时使用整数索引和字符串名称。不允许重复的特征(所选索引和名称之间不能有重叠)。输出向量将首先按给定顺序使用选定的索引对特征排序,然后是选定的名称。
假设有以下 DataFrame:
| userFeatures |
|---|
| [0.0, 10.0, 0.5] |
使用 setIndices(1, 2) 选择后两列:
| userFeatures | features |
|---|---|
| [0.0, 10.0, 0.5] | [10.0, 0.5] |
如果输入属性为 ["f1", "f2", "f3"],使用 setNames("f2", "f3") 选择:
| userFeatures | features |
|---|---|
| [0.0, 10.0, 0.5] | [10.0, 0.5] |
| [“f1”, “f2”, “f3”] | [“f2”, “f3”] |
import java.util.Arrays
import org.apache.spark.ml.attribute.{Attribute, AttributeGroup, NumericAttribute}
import org.apache.spark.ml.feature.VectorSlicer
import org.apache.spark.ml.linalg.Vectors
import org.apache.spark.sql.{Row, SparkSession}
import org.apache.spark.sql.types.StructType
val data = Arrays.asList(
Row(Vectors.sparse(3, Seq((0, -2.0), (1, 2.3)))),
Row(Vectors.dense(-2.0, 2.3, 0.0))
)
val defaultAttr = NumericAttribute.defaultAttr
val attrs = Array("f1", "f2", "f3").map(defaultAttr.withName)
val attrGroup = new AttributeGroup("userFeatures", attrs.asInstanceOf[Array[Attribute]])
val dataset = spark.createDataFrame(data, StructType(Array(attrGroup.toStructField())))
val slicer = new VectorSlicer().setInputCol("userFeatures").setOutputCol("features")
slicer.setIndices(Array(1)).setNames(Array("f3"))
// or slicer.setIndices(Array(1, 2)), or slicer.setNames(Array("f2", "f3"))
val output = slicer.transform(dataset)
output.show(false)
setIndices(Array(1))表示保留索引为1的元素setNames(Array("f3"))表示保留名称为f3的元素setIndices和setNames同时使用时,取并集,保留两者的结果
RFormula
可以将 R 语言风格的公式应用于 DataFrame 数据。
基础操作:
| 符号 | 说明 |
|---|---|
~ | 分离目标变量和项(terms) |
+ | 连接项,"+ 0" 表示移除截距 |
- | 移除一个项,"- 1" 表示移除截距 |
: | 交互项(数值乘法或二值化分类值的交互) |
. | 除目标变量外的所有列 |
假设 a 和 b 是 Double 列:
y ~ a + b→ 模型y ~ w0 + w1*a + w2*b,其中w0是截距,w1、w2是系数y ~ a + b + a:b - 1→ 模型y ~ w1*a + w2*b + w3*a*b,其中w1、w2、w3是系数
RFormula 产生一个特征向量列和一个 Double 或 String 类型的标签列。字符串列会先通过 StringIndexer 转换,然后进行独热编码。
stringOrderType 控制编码方式(假设字符串特征包含 {'b', 'a', 'b', 'a', 'c', 'b'}):
| stringOrderType | StringIndexer 映射到 0 的类别 | RFormula 丢弃的类别 |
|---|---|---|
frequencyDesc | 最频繁类别(b) | 最不频繁类别(c) |
frequencyAsc | 最不频繁类别(c) | 最频繁类别(b) |
alphabetDesc | 字母序最后的类别(c) | 字母序最前的类别(a) |
alphabetAsc | 字母序最前的类别(a) | 字母序最后的类别(c) |
假设有以下 DataFrame:
| id | country | hour | clicked |
|---|---|---|---|
| 7 | ”US” | 18 | 1.0 |
| 8 | ”CA” | 12 | 0.0 |
| 9 | ”NZ” | 15 | 0.0 |
使用公式 clicked ~ country + hour:
| id | country | hour | clicked | features | label |
|---|---|---|---|---|---|
| 7 | ”US” | 18 | 1.0 | [0.0, 0.0, 18.0] | 1.0 |
| 8 | ”CA” | 12 | 0.0 | [0.0, 1.0, 12.0] | 0.0 |
| 9 | ”NZ” | 15 | 0.0 | [1.0, 0.0, 15.0] | 0.0 |
import org.apache.spark.ml.feature.RFormula
val dataset = spark.createDataFrame(Seq(
(7, "US", 18, 1.0),
(8, "CA", 12, 0.0),
(9, "NZ", 15, 0.0)
)).toDF("id", "country", "hour", "clicked")
val formula = new RFormula()
.setFormula("clicked ~ country + hour")
.setFeaturesCol("features")
.setLabelCol("label")
val output = formula.fit(dataset).transform(dataset)
output.select("features", "label").show()
注意:
stringOrderType选项不用于 label 列。当 label 列被索引时,它使用StringIndexer中默认的降频排序。
ChiSqSelector
卡方检验用于测量变量之间的关联性,选择最相关的特征列。
ChiSqSelector 代表卡方特征选择,对具有分类特征的标记数据进行操作,使用独立性的卡平方检验来决定选择哪些特性。支持五种选择方法:
| 方法 | 说明 |
|---|---|
numTopFeatures | 根据卡方检验,选择固定数量的头部特征,类似生成具有最大预测能力的特征 |
percentile | 根据卡方检验,按百分比选择所有特征的一小部分 |
fpr | 根据假阳率(FPR,False Positive Rate),选择 FPR 低于指定值的所有特征。主要用于二分类问题,降低误诊率,常见于医学诊断、信息检索领域 |
fdr | 根据假发现率(FDR,False Discovery Rate),选择错误发现率低于阈值的所有特征。主要用于多重比较问题,常见于统计假设检验、基因表达分析、脑成像研究等领域 |
fwe | 根据 FWER 选择特征,选择 p 值低于阈值的所有特征。阈值按 1/numFeatures 缩放。FWER 衡量多次假设检验中至少有一个错误的概率 |
- 默认选择方法是
numTopFeatures,默认 top features 数量为50- 可以使用
setSelectorType选择选择方法- FPR 和 FDR 主要用于二分类,FPR 关注误诊率,FDR 关注多重比较的错误控制
- FDR 和 FWER 主要用于多重比较问题,FDR 关注虚假发现,FWER 关注整体错误率
假设有以下 DataFrame:
| id | features | clicked |
|---|---|---|
| 7 | [0.0, 0.0, 18.0, 1.0] | 1.0 |
| 8 | [0.0, 1.0, 12.0, 0.0] | 0.0 |
| 9 | [1.0, 0.0, 15.0, 0.1] | 0.0 |
使用 numTopFeatures=1:
| id | features | clicked | selectedFeatures |
|---|---|---|---|
| 7 | [0.0, 0.0, 18.0, 1.0] | 1.0 | [1.0] |
| 8 | [0.0, 1.0, 12.0, 0.0] | 0.0 | [0.0] |
| 9 | [1.0, 0.0, 15.0, 0.1] | 0.0 | [0.1] |
import org.apache.spark.ml.feature.ChiSqSelector
import org.apache.spark.ml.linalg.Vectors
val data = Seq(
(7, Vectors.dense(0.0, 0.0, 18.0, 1.0), 1.0),
(8, Vectors.dense(0.0, 1.0, 12.0, 0.0), 0.0),
(9, Vectors.dense(1.0, 0.0, 15.0, 0.1), 0.0)
)
val df = spark.createDataset(data).toDF("id", "features", "clicked")
val selector = new ChiSqSelector()
.setNumTopFeatures(1)
.setFeaturesCol("features")
.setLabelCol("clicked")
.setOutputCol("selectedFeatures")
val result = selector.fit(df).transform(df)
println(s"ChiSqSelector output with top ${selector.getNumTopFeatures} features selected")
result.show()
setNumTopFeatures(5)选择前 5 个相关特征setPercentile(0.5)选择前 50% 的相关特征setFpr(0.5)选择 FPR 小于 0.5 的相关特征,越小越严格setFdr(0.5)选择 FDR 小于 0.5 的相关特征,越小越严格setFwe(0.5)选择 FWER 小于 0.5 的相关特征,越小越严格
分类与回归
二项逻辑回归(Binomial Logistic Regression)
LogisticRegression 属于监督学习算法,用于解决分类问题,特别是二分类问题。将输入特征的线性组合通过逻辑函数映射到一个介于 0 和 1 之间的值,概率大于等于 0.5 表示分类为正类。
import org.apache.spark.ml.classification.LogisticRegression
// Load training data
val training = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
val lr = new LogisticRegression()
.setMaxIter(10)
.setRegParam(0.3)
.setElasticNetParam(0.8)
// Fit the model
val lrModel = lr.fit(training)
// Print the coefficients and intercept for logistic regression
println(s"Coefficients: ${lrModel.coefficients} Intercept: ${lrModel.intercept}")
// We can also use the multinomial family for binary classification
val mlr = new LogisticRegression()
.setMaxIter(10)
.setRegParam(0.3)
.setElasticNetParam(0.8)
.setFamily("multinomial")
val mlrModel = mlr.fit(training)
// Print the coefficients and intercepts for logistic regression with multinomial family
println(s"Multinomial coefficients: ${mlrModel.coefficientMatrix}")
println(s"Multinomial intercepts: ${mlrModel.interceptVector}")
参数说明
| 参数 | 说明 |
|---|---|
setLabelCol(value: String) | 设置目标列的名称 |
setFeaturesCol(value: String) | 设置特征列的名称 |
setPredictionCol(value: String) | 设置输出预测结果列的名称,一般为 0 或 1 |
setProbabilityCol(value: String) | 设置输出概率列名 |
setRawPredictionCol(value: String) | 设置输出原始预测结果列名 |
setMaxIter(100) | 设置训练时的最大迭代次数 |
setRegParam(value: Double) | 设置正则化参数(lambda),防止过拟合 |
setElasticNetParam(value: Double) | 设置弹性网参数,控制 L1 和 L2 正则化权衡。0.0 为 L2,1.0 为 L1 |
setFamily(value: String) | 设置损失函数类型:auto、binomial(二分类)、multinomial(多分类) |
setThreshold(value: String) | 设置分类阈值,概率值大于此阈值时分类为正类 |
setStandardization(value: String) | 设置是否对特征进行标准化 |
setTol(value: Double) | 设置优化算法的收敛阈值 |
setFitIntercept(value: Boolean) | 设置是否拟合截距 |
setWeightCol(value: String) | 设置样本权重的列名 |
setThresholds(value: Array[Double]) | 设置多分类问题的分类阈值 |
setLowerBoundsOnIntercepts(value: Array[Double]) | 限制截距项的最小值 |
setUpperBoundsOnIntercepts(value: Array[Double]) | 限制截距项的最大值 |
setLowerBoundsOnCoefficients(value: Array[Double]) | 限制多个系数的最小值 |
setUpperBoundsOnCoefficients(value: Array[Double]) | 限制多个系数的最大值 |
setAggregationDepth(value: Int) | 计算梯度和 Hessian 矩阵的聚合深度 |
setAggregationDepth:聚合深度越大,将在更多数据分区之间执行数据聚合操作,以提高算法性能。内存有限时建议较低聚合深度;数据分布极端不平衡时建议较高聚合深度;模型参数数量较多时需要更高聚合深度。
模型摘要(Binary Summary)
LogisticRegressionTrainingSummary 提供 LogisticRegressionModel 的摘要。在二进制分类的情况下,还可以使用某些额外的度量(如 ROC 曲线)。可以通过 binarySummary 方法访问二进制摘要。
import org.apache.spark.ml.classification.LogisticRegression
// Extract the summary from the returned LogisticRegressionModel instance
val trainingSummary = lrModel.binarySummary
// Obtain the objective per iteration.
val objectiveHistory = trainingSummary.objectiveHistory
println("objectiveHistory:")
objectiveHistory.foreach(loss => println(loss))
// Obtain the receiver-operating characteristic as a dataframe and areaUnderROC.
val roc = trainingSummary.roc
roc.show()
println(s"areaUnderROC: ${trainingSummary.areaUnderROC}")
// Set the model threshold to maximize F-Measure
val fMeasure = trainingSummary.fMeasureByThreshold
val maxFMeasure = fMeasure.select(max("F-Measure")).head().getDouble(0)
val bestThreshold = fMeasure.where($"F-Measure" === maxFMeasure)
.select("threshold").head().getDouble(0)
lrModel.setThreshold(bestThreshold)
模型信息
以 lrModel 为例:
| 属性/方法 | 说明 |
|---|---|
coefficients | 二分类模型系数 |
coefficientMatrix | 多分类模型系数 |
intercept | 二分类模型截距项 |
interceptVector | 多分类模型截距项 |
binarySummary | 二分类问题的模型摘要 |
binarySummary 可获取的信息:
| 属性/方法 | 说明 |
|---|---|
areaUnderROC | ROC 曲线下面积,0 到 1 之间,越接近 1 越好 |
roc | ROC 曲线,包含不同阈值下的真正例比率和假正例比率的 DataFrame |
pr | PR(Precision-Recall)曲线,包含不同阈值下的精确率和召回率的 DataFrame |
fMeasureByThreshold | 根据不同阈值获取 F1 分数 |
precisionByThreshold | 根据不同阈值获取精确率 |
recallByThreshold | 根据不同阈值获取召回率 |
defaultThreshold | 获取模型的默认分类阈值(通常为 0.5) |
prThresholds | 获取 PR 曲线中使用的不同阈值数组 |
rocThresholds | 获取 ROC 曲线中使用的不同阈值数组 |
objectiveHistory | 每次迭代中损失函数值的变化情况 |
多项逻辑回归(Multinomial Logistic Regression)
通过多项式 logistic 回归(softmax)支持多类分类。在多项式 logistic 回归中,算法产生 K 组系数,或 K × J 维矩阵,其中 K 是结果类的数目,J 是特征的数目。
结果类别 k ∈ 1,2,…,K 的条件概率使用 softmax 函数建模:
P(Y=k|X,βk,β0k) = e^(βk·X+β0k) / Σ e^(βk′·X+β0k′)
最小化加权负对数似然,使用多项式响应模型,带弹性网惩罚控制过拟合:
min[β,β0] -[Σ wi·logP(Y=yi|xi)] + λ[(1/2)(1-α)||β||₂² + α||β||₁]
注意:当模型为二分类时,系数保存在
coefficients属性中,截距项保存在intercept属性中;当模型为多分类时,系数保存在coefficientMatrix属性中,截距项保存在interceptVector属性中。
import org.apache.spark.ml.classification.LogisticRegression
// Load training data
val training = spark
.read
.format("libsvm")
.load("data/mllib/sample_multiclass_classification_data.txt")
val lr = new LogisticRegression()
.setMaxIter(10)
.setRegParam(0.3)
.setElasticNetParam(0.8)
// Fit the model
val lrModel = lr.fit(training)
// Print the coefficients and intercept for multinomial logistic regression
println(s"Coefficients: \n${lrModel.coefficientMatrix}")
println(s"Intercepts: \n${lrModel.interceptVector}")
val trainingSummary = lrModel.summary
// Obtain the objective per iteration
val objectiveHistory = trainingSummary.objectiveHistory
println("objectiveHistory:")
objectiveHistory.foreach(println)
// for multiclass, we can inspect metrics on a per-label basis
println("False positive rate by label:")
trainingSummary.falsePositiveRateByLabel.zipWithIndex.foreach { case (rate, label) =>
println(s"label $label: $rate")
}
println("True positive rate by label:")
trainingSummary.truePositiveRateByLabel.zipWithIndex.foreach { case (rate, label) =>
println(s"label $label: $rate")
}
println("Precision by label:")
trainingSummary.precisionByLabel.zipWithIndex.foreach { case (prec, label) =>
println(s"label $label: $prec")
}
println("Recall by label:")
trainingSummary.recallByLabel.zipWithIndex.foreach { case (rec, label) =>
println(s"label $label: $rec")
}
println("F-measure by label:")
trainingSummary.fMeasureByLabel.zipWithIndex.foreach { case (f, label) =>
println(s"label $label: $f")
}
val accuracy = trainingSummary.accuracy
val falsePositiveRate = trainingSummary.weightedFalsePositiveRate
val truePositiveRate = trainingSummary.weightedTruePositiveRate
val fMeasure = trainingSummary.weightedFMeasure
val precision = trainingSummary.weightedPrecision
val recall = trainingSummary.weightedRecall
println(s"Accuracy: $accuracy\nFPR: $falsePositiveRate\nTPR: $truePositiveRate\n" +
s"F-measure: $fMeasure\nPrecision: $precision\nRecall: $recall")
多分类模型摘要
以 lrModel 为例,summary 属性(适用于二分类和多分类问题)可获取以下信息:
| 属性/方法 | 说明 |
|---|---|
totalIteration: Int | 模型训练的总迭代次数 |
objectiveHistory: Array[Double] | 每次迭代的损失函数值数组,用于绘制收敛情况 |
precisionByLabel: Array[Double] | 每个类别的精确率 |
recallByLabel: Array[Double] | 每个类别的召回率 |
fMeasureByLabel: Array[Double] | 每个类别的 F1 Score |
weightedFalsePositiveRate: Double | 加权的假阳性率 |
weightedTruePositiveRate: Double | 加权的真阳性率 |
weightedFMeasure: Double | 加权的 F1 Score |
weightedPrecision: Double | 加权的精确率 |
weightedRecall: Double | 加权的召回率 |
accuracy: Double | 模型的准确率 |
confusionMatrix: Matrix | 混淆矩阵,包含实际类别和预测类别之间的统计信息 |
决策树分类器(Decision Tree Classifier)
决策树分类器用于解决分类问题,目标变量是离散的类别,也可以用于多分类问题。
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.classification.DecisionTreeClassificationModel
import org.apache.spark.ml.classification.DecisionTreeClassifier
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
import org.apache.spark.ml.feature.{IndexToString, StringIndexer, VectorIndexer}
// Load the data stored in LIBSVM format as a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Index labels, adding metadata to the label column.
// Fit on whole dataset to include all labels in index.
val labelIndexer = new StringIndexer()
.setInputCol("label")
.setOutputCol("indexedLabel")
.fit(data)
// Automatically identify categorical features, and index them.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4) // features with > 4 distinct values are treated as continuous.
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a DecisionTree model.
val dt = new DecisionTreeClassifier()
.setLabelCol("indexedLabel")
.setFeaturesCol("indexedFeatures")
// Convert indexed labels back to original labels.
val labelConverter = new IndexToString()
.setInputCol("prediction")
.setOutputCol("predictedLabel")
.setLabels(labelIndexer.labels)
// Chain indexers and tree in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(labelIndexer, featureIndexer, dt, labelConverter))
// Train model. This also runs the indexers.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("predictedLabel", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new MulticlassClassificationEvaluator()
.setLabelCol("indexedLabel")
.setPredictionCol("prediction")
.setMetricName("accuracy")
val accuracy = evaluator.evaluate(predictions)
println(s"Test Error = ${(1.0 - accuracy)}")
val treeModel = model.stages(2).asInstanceOf[DecisionTreeClassificationModel]
println(s"Learned classification tree model:\n ${treeModel.toDebugString}")
参数说明
| 参数 | 说明 |
|---|---|
setLabelCol(value: String) | 设置标签列的列名 |
setFeaturesCol(value: String) | 设置特征列的列名 |
setPredictionCol(value: String) | 设置模型预测结果的列名 |
setImpurity(value: String) | 设置分割特征的不纯度度量:gini(基尼不纯度)、entropy(信息增益) |
setMaxDepth(value: Int) | 设置树的最大深度,防止过拟合 |
setMaxBins(value: Int) | 设置分割特征的最大分箱数量,加速特征分割计算 |
setMinInstancesPerNode(value: Int) | 设置每个节点分割前所需的最小实例数 |
setMinInfoGain(value: Double) | 设置每次分割所需的最小信息增益 |
setSeed(seed: Long) | 设置随机种子 |
setMaxMemoryInMB(value: Int) | 设置分割特征时的内存限制 |
setImpurity:gini不纯度值范围 0 到 0.5,0 表示完全纯净,0.5 表示完全不纯净;entropy信息增益值范围 0 到 1,0 表示完全纯净,1 表示完全不纯净。特征分割概念:构建决策树时,算法会选择一个特征和相应的阈值,将数据集分成两个子集。分箱可以将连续特征值划分为离散区间,加速特征分割,同时降低模型复杂度,减少过拟合风险。
随机森林分类器(Random Forest Classifier)
随机森林是一种 Bagging 集成学习方法,通过组合多个决策树来进行分类,提高模型的性能和鲁棒性。主要特点是每个决策树使用特征集合的一个随机子集进行训练,最后通过多数投票或平均值来确定最终分类结果。
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.classification.{RandomForestClassificationModel, RandomForestClassifier}
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
import org.apache.spark.ml.feature.{IndexToString, StringIndexer, VectorIndexer}
// Load and parse the data file, converting it to a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Index labels, adding metadata to the label column.
val labelIndexer = new StringIndexer()
.setInputCol("label")
.setOutputCol("indexedLabel")
.fit(data)
// Automatically identify categorical features, and index them.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4)
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a RandomForest model.
val rf = new RandomForestClassifier()
.setLabelCol("indexedLabel")
.setFeaturesCol("indexedFeatures")
.setNumTrees(10)
// Convert indexed labels back to original labels.
val labelConverter = new IndexToString()
.setInputCol("prediction")
.setOutputCol("predictedLabel")
.setLabels(labelIndexer.labels)
// Chain indexers and forest in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(labelIndexer, featureIndexer, rf, labelConverter))
// Train model. This also runs the indexers.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("predictedLabel", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new MulticlassClassificationEvaluator()
.setLabelCol("indexedLabel")
.setPredictionCol("prediction")
.setMetricName("accuracy")
val accuracy = evaluator.evaluate(predictions)
println(s"Test Error = ${(1.0 - accuracy)}")
val rfModel = model.stages(2).asInstanceOf[RandomForestClassificationModel]
println(s"Learned classification forest model:\n ${rfModel.toDebugString}")
参数说明
| 参数 | 说明 |
|---|---|
setNumTrees(value: Int) | 设置随机森林中树的数量 |
setMaxDepth(value: Int) | 设置每棵树的最大深度,防止过拟合 |
setMinInstancesPerNode(value: Int) | 设置每个节点分割前所需的最小实例数 |
setFeatureSubsetStrategy(value: String) | 设置特征选择策略:auto、all、sqrt、onethird |
setImpurity(value: String) | 设置不纯度度量:gini、entropy |
setSeed(value: Long) | 设置随机种子 |
setMaxBins(value: Int) | 设置最大分箱数量 |
setCacheNodeIds(value: Boolean) | 设置是否缓存节点 ID,加速特征重要性计算 |
setFeatureSubsetStrategy:auto根据问题类型自动选择;all使用所有特征;sqrt使用特征数平方根的子集;onethird使用特征数三分之一的子集。特征重要性:衡量每个特征对模型预测的贡献程度。随机森林中的特征重要性通常基于每个特征在构建树时的节点分割中的贡献度来计算。
setCacheNodeIds缓存节点 ID 可以避免重复遍历,节省计算资源和时间。
梯度提升树分类器(GBT Classifier)
GBT 是一种 Boosting 集成学习方法,通过组合多个决策树来进行分类,提高模型的性能和鲁棒性。主要特点是顺序训练多个基学习器,每个基学习器都试图纠正前一个学习器的错误。
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.classification.{GBTClassificationModel, GBTClassifier}
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
import org.apache.spark.ml.feature.{IndexToString, StringIndexer, VectorIndexer}
// Load and parse the data file, converting it to a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Index labels, adding metadata to the label column.
val labelIndexer = new StringIndexer()
.setInputCol("label")
.setOutputCol("indexedLabel")
.fit(data)
// Automatically identify categorical features, and index them.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4)
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a GBT model.
val gbt = new GBTClassifier()
.setLabelCol("indexedLabel")
.setFeaturesCol("indexedFeatures")
.setMaxIter(10)
.setFeatureSubsetStrategy("auto")
// Convert indexed labels back to original labels.
val labelConverter = new IndexToString()
.setInputCol("prediction")
.setOutputCol("predictedLabel")
.setLabels(labelIndexer.labels)
// Chain indexers and GBT in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(labelIndexer, featureIndexer, gbt, labelConverter))
// Train model. This also runs the indexers.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("predictedLabel", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new MulticlassClassificationEvaluator()
.setLabelCol("indexedLabel")
.setPredictionCol("prediction")
.setMetricName("accuracy")
val accuracy = evaluator.evaluate(predictions)
println(s"Test Error = ${1.0 - accuracy}")
val gbtModel = model.stages(2).asInstanceOf[GBTClassificationModel]
println(s"Learned classification GBT model:\n ${gbtModel.toDebugString}")
参数说明
| 参数 | 说明 |
|---|---|
setMaxIter(value: Int) | 设置训练迭代次数,即构建树的数量 |
setMaxDepth(value: Int) | 设置每棵树的最大深度 |
setMaxBins(value: Int) | 设置最大分箱数量 |
setStepSize(stepSize: Double) | 设置学习率,影响梯度提升步长 |
setSeed(seed: Long) | 设置随机种子 |
setFeatureSubsetStrategy(value: String) | 设置特征选择策略:all、sqrt、onethird |
setImpurity(value: String) | 设置不纯度度量:gini、entropy |
Bagging vs Boosting 对比
| Bagging(随机森林) | Boosting(梯度提升树) | |
|---|---|---|
| 训练方式 | 并行训练多棵树 | 顺序训练,每棵树纠正前一棵的错误 |
| 样本权重 | 无权重区分 | 错误分类样本获得更高权重 |
| 特征选择 | 随机选择特征子集 | 支持特征子集选择 |
| 多样性 | 通过随机特征子集降低相关性 | 通过权重调整增加多样性 |
| 可解释性 | 较低(多棵树组成) | 较低(多棵树组成) |
多层感知机分类器(MLPC)
MLP 是一种深度神经网络,通常包括多个隐藏层,用于解决分类问题。MLP 是一种前馈神经网络,包括输入层、多个隐藏层和输出层。每个神经元与前一层的神经元相连,通过权重和激活函数进行计算。
中间层使用 sigmoid(logistic)函数:
f(zi) = 1 / (1 + e^(-zi))
输出层使用 softmax 函数:
f(zi) = e^zi / Σ e^zk
import org.apache.spark.ml.classification.MultilayerPerceptronClassifier
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
// Load the data stored in LIBSVM format as a DataFrame.
val data = spark.read.format("libsvm")
.load("data/mllib/sample_multiclass_classification_data.txt")
// Split the data into train and test
val splits = data.randomSplit(Array(0.6, 0.4), seed = 1234L)
val train = splits(0)
val test = splits(1)
// specify layers for the neural network:
// input layer of size 4 (features), two intermediate of size 5 and 4
// and output of size 3 (classes)
val layers = Array[Int](4, 5, 4, 3)
// create the trainer and set its parameters
val trainer = new MultilayerPerceptronClassifier()
.setLayers(layers)
.setBlockSize(128)
.setSeed(1234L)
.setMaxIter(100)
// train the model
val model = trainer.fit(train)
// compute accuracy on the test set
val result = model.transform(test)
val predictionAndLabels = result.select("prediction", "label")
val evaluator = new MulticlassClassificationEvaluator()
.setMetricName("accuracy")
println(s"Test set accuracy = ${evaluator.evaluate(predictionAndLabels)}")
参数说明
| 参数 | 说明 |
|---|---|
setLayers(value: Array[Int]) | 设置网络层次结构,首元素为输入维度,尾元素为输出类别数,中间为隐藏层大小 |
setBlockSize(value: Int) | 设置并行处理的块大小,较大块提高性能但增加内存需求 |
setSeed(seed: Long) | 设置随机种子 |
setMaxIter(value: Int) | 设置最大迭代次数 |
setTol(tol: Double) | 设置迭代收敛阈值 |
setStepSize(stepSize: Double) | 设置学习率,影响梯度下降步长 |
setSolver(value: String) | 设置优化算法:l-bfgs、gd |
setLayers:以[4, 5, 4, 3]为例,第一个参数4为输入特征维度,最后一个参数3为输出类别数量,中间5和4为隐藏层大小。如果输入为 OneHot 编码的特征列,第一个参数即为稀疏向量的维度。数据分批训练:在大规模数据集上分批训练通常不会显著影响模型质量,前提是每个数据块包含足够的随机性、块之间正确传递信息、验证和测试数据处理方式一致。
线性支持向量机(Linear SVC)
支持向量机分类器用于解决分类问题,属于监督学习算法。模型试图找到一个超平面,将不同类别的数据分隔开。
LinearSVC 在 Spark ML 中使用 OWLQN 优化器优化 Hinge Loss:
超平面:w^T - b = 0
分类决策:f(x) = sign(w^T - b)
损失函数:L(y, f(x)) = max(0, 1 - y(w^T - b))
目标函数:min 1/2 ∥w∥² + C * Σ L(yi, f(xi))
import org.apache.spark.ml.classification.LinearSVC
// Load training data
val training = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
val lsvc = new LinearSVC()
.setMaxIter(10)
.setRegParam(0.1)
// Fit the model
val lsvcModel = lsvc.fit(training)
// Print the coefficients and intercept for linear svc
println(s"Coefficients: ${lsvcModel.coefficients} Intercept: ${lsvcModel.intercept}")
参数说明
| 参数 | 说明 |
|---|---|
setMaxIter(value: Int) | 设置最大迭代次数 |
setRegParam(value: Double) | 设置正则化参数(lambda),防止过拟合 |
setTol(tol: Double) | 设置收敛阈值 |
setStandardization(value: Boolean) | 设置是否标准化特征 |
setFitIntercept(value: Boolean) | 设置是否拟合截距 |
setThreshold(value: Double) | 设置分类阈值,默认为 0.0 |
setElasticNetParam(value: Double) | 设置弹性网参数 |
One-vs-Rest 分类器
OneVsRest 是多分类问题中的一种策略,允许使用多个二元分类器处理多类别问题。对于具有 K 个类别的问题,训练 K 个二元分类器,每个负责将一个类别与其他所有类别区分。预测时选择具有最高分数的类别。
import org.apache.spark.ml.classification.{LogisticRegression, OneVsRest}
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
// load data file.
val inputData = spark.read.format("libsvm")
.load("data/mllib/sample_multiclass_classification_data.txt")
// generate the train/test split.
val Array(train, test) = inputData.randomSplit(Array(0.8, 0.2))
// instantiate the base classifier
val classifier = new LogisticRegression()
.setMaxIter(10)
.setTol(1E-6)
.setFitIntercept(true)
// instantiate the One Vs Rest Classifier.
val ovr = new OneVsRest().setClassifier(classifier)
// train the multiclass model.
val ovrModel = ovr.fit(train)
// score the model on test data.
val predictions = ovrModel.transform(test)
// obtain evaluator.
val evaluator = new MulticlassClassificationEvaluator()
.setMetricName("accuracy")
// compute the classification error on test data.
val accuracy = evaluator.evaluate(predictions)
println(s"Test Error = ${1 - accuracy}")
参数说明
| 参数 | 说明 |
|---|---|
setClassifier(value) | 设置基础二元分类器(如 LogisticRegression) |
setMetricName(value: String) | 设置评估度量:accuracy、f1、weightedPrecision、weightedRecall、weightedF1、logloss |
setLabelCol(value: String) | 设置标签列 |
setFeaturesCol(value: String) | 设置特征列 |
setPredictionCol(value: String) | 设置预测结果列 |
setWeightCol(value: String) | 设置样本权重列 |
setThresholds(value: Array[Double]) | 设置多分类阈值数组 |
setRawPredictionCol(value: String) | 设置原始预测结果列(Vector 类型) |
setParallelism(value: Int) | 设置并行度,即同时训练的二元分类器数量 |
推荐:使用
setMetricName设置评估度量,用于比较不同模型或调整模型参数。
朴素贝叶斯分类器(Naive Bayes)
朴素贝叶斯基于特征之间的条件独立假设(“朴素”的含义),利用贝叶斯定理计算给定特征条件下的类别概率。
MLlib 支持两种模型类型:
- Multinomial(多项式):特征表示词频(默认)
- Bernoulli(伯努利):特征表示是否出现(0 或 1)
特征值必须为非负数。可以使用平滑参数 λ(默认为 1.0)进行加法平滑处理。
import org.apache.spark.ml.classification.NaiveBayes
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
// Load the data stored in LIBSVM format as a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Split the data into training and test sets (30% held out for testing)
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3), seed = 1234L)
// Train a NaiveBayes model.
val model = new NaiveBayes()
.fit(trainingData)
// Select example rows to display.
val predictions = model.transform(testData)
predictions.show()
// Select (prediction, true label) and compute test error
val evaluator = new MulticlassClassificationEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("accuracy")
val accuracy = evaluator.evaluate(predictions)
println(s"Test set accuracy = $accuracy")
参数说明
| 参数 | 说明 |
|---|---|
setLabelCol(value: String) | 设置标签列 |
setFeaturesCol(value: String) | 设置特征列 |
setPredictionCol(value: String) | 设置预测结果列 |
setSmoothing(value: Double) | 设置平滑参数,处理零频率问题(拉普拉斯平滑) |
setModelType(value: String) | 设置模型类型:multinomial(多项式)、bernoulli(伯努利) |
setWeightCol(value: String) | 设置样本权重列 |
setThreshold(value: Double) | 设置二元分类阈值 |
setThresholds(value: Array[Double]) | 设置多分类阈值数组 |
setRawPredictionCol(value: String) | 设置原始预测结果列 |
setProbabilityCol(value: String) | 设置概率列 |
setAggregationDepth(value: Int) | 设置聚合深度 |
模型信息
| 属性 | 说明 |
|---|---|
pi: DenseVector | 每个类别的先验概率向量 |
theta: Array[DenseVector] | 每个类别下每个特征的条件概率矩阵数组 |
回归
线性回归(Linear Regression)
建立线性关系模型,用于解决回归问题。与分类的离散预测值不同,回归的目标是预测连续数值的输出,通过学习权重和偏置来拟合数据,最小化损失函数(通常是均方误差)。
模型公式:y = w1*x1 + w2*x2 + ... + wn*xn + b
import org.apache.spark.ml.regression.LinearRegression
// Load training data
val training = spark.read.format("libsvm")
.load("data/mllib/sample_linear_regression_data.txt")
val lr = new LinearRegression()
.setMaxIter(10)
.setRegParam(0.3)
.setElasticNetParam(0.8)
// Fit the model
val lrModel = lr.fit(training)
// Print the coefficients and intercept for linear regression
println(s"Coefficients: ${lrModel.coefficients} Intercept: ${lrModel.intercept}")
// Summarize the model over the training set and print out some metrics
val trainingSummary = lrModel.summary
println(s"numIterations: ${trainingSummary.totalIterations}")
println(s"objectiveHistory: [${trainingSummary.objectiveHistory.mkString(",")}]")
trainingSummary.residuals.show()
println(s"RMSE: ${trainingSummary.rootMeanSquaredError}")
println(s"r2: ${trainingSummary.r2}")
参数说明
| 参数 | 说明 |
|---|---|
setMaxIter(100) | 设置最大迭代次数 |
setRegParam(value: Double) | 设置正则化参数(lambda) |
setElasticNetParam(value: Double) | 设置弹性网参数,0.0 为 L2,1.0 为 L1 |
setStandardization(value: Boolean) | 设置是否标准化特征 |
setTol(value: Double) | 设置收敛阈值 |
setFitIntercept(value: Boolean) | 设置是否拟合截距 |
setWeightCol(value: String) | 设置样本权重列 |
setAggregationDepth(value: Int) | 设置梯度/Hessian 矩阵聚合深度 |
广义线性回归(Generalized Linear Regression)
适用于处理各种线性回归问题,包括线性回归、泊松回归、二项分布回归等。设计灵感来自广义线性模型(GLM)理论,用于解决回归问题,目标是预测连续数值输出。
模型公式:g(E(y)) = β0 + β1*x1 + β2*x2 + ... + βn*xn
其中 g 是链接函数,E(y) 是期望输出。
注意:Spark 目前最高仅支持 4096 个特征通过
GeneralizedLinearRegression接口。
可用分布族
| 分布族 | 返回类型 | 支持的链接函数 |
|---|---|---|
| Gaussian | Continuous | Identity*、Log、Inverse |
| Binomial | Binary | Logit*、Probit、CLogLog |
| Poisson | Count | Log*、Identity、Sqrt |
| Gamma | Continuous | Inverse*、Identity、Log |
| Tweedie | Zero-inflated continuous | Power link function |
*带 * 的为 Canonical Link(标准链接函数)
import org.apache.spark.ml.regression.GeneralizedLinearRegression
// Load training data
val dataset = spark.read.format("libsvm")
.load("data/mllib/sample_linear_regression_data.txt")
val glr = new GeneralizedLinearRegression()
.setFamily("gaussian")
.setLink("identity")
.setMaxIter(10)
.setRegParam(0.3)
// Fit the model
val model = glr.fit(dataset)
// Print the coefficients and intercept for generalized linear regression model
println(s"Coefficients: ${model.coefficients}")
println(s"Intercept: ${model.intercept}")
// Summarize the model over the training set and print out some metrics
val summary = model.summary
println(s"Coefficient Standard Errors: ${summary.coefficientStandardErrors.mkString(",")}")
println(s"T Values: ${summary.tValues.mkString(",")}")
println(s"P Values: ${summary.pValues.mkString(",")}")
println(s"Dispersion: ${summary.dispersion}")
println(s"Null Deviance: ${summary.nullDeviance}")
println(s"Residual Degree Of Freedom Null: ${summary.residualDegreeOfFreedomNull}")
println(s"Deviance: ${summary.deviance}")
println(s"Residual Degree Of Freedom: ${summary.residualDegreeOfFreedom}")
println(s"AIC: ${summary.aic}")
println("Deviance Residuals: ")
summary.residuals().show()
参数说明
| 参数 | 说明 |
|---|---|
setMaxIter(100) | 设置最大迭代次数 |
setRegParam(value: Double) | 设置正则化参数(lambda) |
setLink(value: String) | 设置链接函数:identity、log、inverse、logit、probit |
setFamily(value: String) | 设置分布族:gaussian、binomial、poisson、gamma、tweedie |
setVariancePower(value: Double) | 设置方差函数的幂(用于 Tweedie 分布) |
setStandardization(value: Boolean) | 设置是否标准化特征 |
setTol(value: Double) | 设置收敛阈值 |
setFitIntercept(value: Boolean) | 设置是否拟合截距 |
setFamily分布族选择指南:
gaussian:连续数值回归问题,假设误差服从正态分布binomial:二值回归问题(逻辑回归)poisson:计数型回归问题,计数数据、小概率事件发生率估计gamma:正数数据回归、偏斜数据回归、生存分析tweedie:更灵活的分布,alpha=1等价于泊松,alpha=2等价于伽马
setLink链接函数选择指南:
identity:恒等函数,等同于线性回归log:对数函数,适用于泊松回归inverse:倒数函数,适用于逆高斯回归logit:逻辑函数,适用于二元分布问题probit:正态分布累积分布函数,适用于正态分布回归
setVariancePower:调整方差的幂次数p。p=0方差与均值无关(正态分布),p=1方差与均值成正比(泊松分布),p≠0 且 p≠1用于非线性方差结构。
决策树回归(Decision Tree Regression)
决策树回归用于解决回归问题,预测连续数值输出。
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.feature.VectorIndexer
import org.apache.spark.ml.regression.DecisionTreeRegressionModel
import org.apache.spark.ml.regression.DecisionTreeRegressor
// Load the data stored in LIBSVM format as a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Automatically identify categorical features, and index them.
// Here, we treat features with > 4 distinct values as continuous.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4)
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a DecisionTree model.
val dt = new DecisionTreeRegressor()
.setLabelCol("label")
.setFeaturesCol("indexedFeatures")
// Chain indexer and tree in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(featureIndexer, dt))
// Train model. This also runs the indexer.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("prediction", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new RegressionEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("rmse")
val rmse = evaluator.evaluate(predictions)
println(s"Root Mean Squared Error (RMSE) on test data = $rmse")
val treeModel = model.stages(1).asInstanceOf[DecisionTreeRegressionModel]
println(s"Learned regression tree model:\n ${treeModel.toDebugString}")
setter 方法参考决策树分类器。额外参数:
setCheckpointInterval(checkpointInterval):设置检查点间隔,提高容错性setLeafCol(leafCol):设置叶子节点的列名
随机森林回归(Random Forest Regression)
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.feature.VectorIndexer
import org.apache.spark.ml.regression.{RandomForestRegressionModel, RandomForestRegressor}
// Load and parse the data file, converting it to a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Automatically identify categorical features, and index them.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4)
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a RandomForest model.
val rf = new RandomForestRegressor()
.setLabelCol("label")
.setFeaturesCol("indexedFeatures")
// Chain indexer and forest in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(featureIndexer, rf))
// Train model. This also runs the indexer.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("prediction", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new RegressionEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("rmse")
val rmse = evaluator.evaluate(predictions)
println(s"Root Mean Squared Error (RMSE) on test data = $rmse")
val rfModel = model.stages(1).asInstanceOf[RandomForestRegressionModel]
println(s"Learned regression forest model:\n ${rfModel.toDebugString}")
梯度提升树回归(GBT Regression)
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.feature.VectorIndexer
import org.apache.spark.ml.regression.{GBTRegressionModel, GBTRegressor}
// Load and parse the data file, converting it to a DataFrame.
val data = spark.read.format("libsvm").load("data/mllib/sample_libsvm_data.txt")
// Automatically identify categorical features, and index them.
val featureIndexer = new VectorIndexer()
.setInputCol("features")
.setOutputCol("indexedFeatures")
.setMaxCategories(4)
.fit(data)
// Split the data into training and test sets (30% held out for testing).
val Array(trainingData, testData) = data.randomSplit(Array(0.7, 0.3))
// Train a GBT model.
val gbt = new GBTRegressor()
.setLabelCol("label")
.setFeaturesCol("indexedFeatures")
.setMaxIter(10)
// Chain indexer and GBT in a Pipeline.
val pipeline = new Pipeline()
.setStages(Array(featureIndexer, gbt))
// Train model. This also runs the indexer.
val model = pipeline.fit(trainingData)
// Make predictions.
val predictions = model.transform(testData)
// Select example rows to display.
predictions.select("prediction", "label", "features").show(5)
// Select (prediction, true label) and compute test error.
val evaluator = new RegressionEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("rmse")
val rmse = evaluator.evaluate(predictions)
println(s"Root Mean Squared Error (RMSE) on test data = $rmse")
val gbtModel = model.stages(1).asInstanceOf[GBTRegressionModel]
println(s"Learned regression GBT model:\n ${gbtModel.toDebugString}")
注意:对于示例数据集,
GBTRegressor实际上只需要 1 次迭代,但一般情况下并非如此。
生存回归(Survival Regression / AFT)
在 spark.ml 中,实现了加速失效时间(AFT)模型,这是一个用于截尾数据的参数化生存回归模型。它描述了生存时间对数的模型,因此常被称为生存分析的对数线性模型。
与同样用于此目的的 Proportional Hazards 模型不同,AFT 模型更容易并行化,因为每个实例独立地对目标函数做贡献。
给定协变量 x',对于样本 i = 1, …, n 的随机生存时间 ti,似然函数为:
L(β,σ) = Π [1/σ * f0((log(ti)-x'β)/σ)]^δi * S0((log(ti)-x'β)/σ)^(1-δi)
其中 δi 是事件发生的指示器。使用 ϵi = (log(ti)-x'β)/σ,对数似然函数为:
ι(β,σ) = Σ [-δi*log(σ) + δi*log(f0(ϵi)) + (1-δi)*log(S0(ϵi))]
最常用的 AFT 模型基于生存时间的 Weibull 分布,其中:
S0(ϵi) = exp(-e^ϵi)
f0(ϵi) = e^ϵi * exp(-e^ϵi)
AFT 模型可表示为凸优化问题,底层使用 L-BFGS 优化算法。
import org.apache.spark.ml.linalg.Vectors
import org.apache.spark.ml.regression.AFTSurvivalRegression
val training = spark.createDataFrame(Seq(
(1.218, 1.0, Vectors.dense(1.560, -0.605)),
(2.949, 0.0, Vectors.dense(0.346, 2.158)),
(3.627, 0.0, Vectors.dense(1.380, 0.231)),
(0.273, 1.0, Vectors.dense(0.520, 1.151)),
(4.199, 0.0, Vectors.dense(0.795, -0.226))
)).toDF("label", "censor", "features")
val quantileProbabilities = Array(0.3, 0.6)
val aft = new AFTSurvivalRegression()
.setQuantileProbabilities(quantileProbabilities)
.setQuantilesCol("quantiles")
val model = aft.fit(training)
// Print the coefficients, intercept and scale parameter for AFT survival regression
println(s"Coefficients: ${model.coefficients}")
println(s"Intercept: ${model.intercept}")
println(s"Scale: ${model.scale}")
model.transform(training).show(false)
等保回归(Isotonic Regression)
等保回归属于回归算法家族。给定一组实数 Y = y1,y2,...,yn 表示观测响应值,以及 X = x1,x2,...,xn 表示待拟合的未知响应值,需要找到函数最小化:
f(x) = Σ wi(yi - xi)²
并满足完全排序条件 x1 ≤ x2 ≤ ... ≤ xn,其中 wi 是正权重。
结果函数称为等保回归,可视为在排序约束下的最小二乘问题,即最适合原始数据点的单调函数。
Spark 实现了一个 Pool Adjacent Violators (PAV) 算法的并行化版本。isotonic 参数默认为 true,指定等保回归是单调递增(true)还是单调递减(false)。
预测规则:
- 如果预测输入恰好等于训练特征,返回关联预测
- 如果预测输入低于/高于所有训练特征,返回最低/最高特征的预测
- 如果预测输入落在两个训练特征之间,使用分段线性函数插值计算
import org.apache.spark.ml.regression.IsotonicRegression
// Loads data.
val dataset = spark.read.format("libsvm")
.load("data/mllib/sample_isotonic_regression_libsvm_data.txt")
// Trains an isotonic regression model.
val ir = new IsotonicRegression()
val model = ir.fit(dataset)
println(s"Boundaries in increasing order: ${model.boundaries}\n")
println(s"Predictions associated with the boundaries: ${model.predictions}\n")
// Makes predictions.
model.transform(dataset).show()
参数说明
| 参数 | 说明 |
|---|---|
setFeaturesCol(value: String) | 设置特征列名称 |
setFeaturesIndex(value: String) | 设置特征索引列名称(多特征时指定使用哪一列) |
setLabelCol(value: String) | 设置目标变量列名 |
setPredictionCol(value: String) | 设置预测结果列名 |
setWeightCol(value: String) | 设置样本权重列 |
setIsotonic(value: Boolean) | 设置是否启用等保回归(单调性),默认 true |
线性方法与弹性网(Elastic Net)
弹性网是 L1 和 L2 正则化的混合,数学上定义为两者的凸组合:
α(λ∥w∥₁) + (1-α)(λ/2 * ∥w∥₂²), α∈[0,1], λ≥0
- 当
α = 1:等价于 Lasso 模型(仅 L1) - 当
α = 0:等价于 Ridge 回归模型(仅 L2)
Spark MLlib 使用 Pipelines API 实现了带弹性网正则化的线性回归和逻辑回归。
决策树详解
决策树及其集成方法是分类和回归任务中最流行的方法之一。决策树被广泛使用,因为它们易于解释、能处理分类特征、扩展到多分类设置、不需要特征缩放、能够捕捉非线性和特征交互。
输入输出列
输入列
| 参数名 | 类型 | 默认值 | 说明 |
|---|---|---|---|
labelCol | Double | "label" | 要预测的标签 |
featuresCol | Vector | "features" | 特征向量 |
输出列
| 参数名 | 类型 | 默认值 | 说明 | 备注 |
|---|---|---|---|---|
predictionCol | Double | "prediction" | 预测标签 | |
rawPredictionCol | Vector | "rawPrediction" | 长度等于类别数的向量 | 仅分类 |
probabilityCol | Vector | "probability" | 归一化为多项分布的 rawPrediction | 仅分类 |
varianceCol | Double | 预测的有偏样本方差 | 仅回归 |
与原始 MLlib 决策树 API 的主要区别:
- 支持 ML Pipelines
- 区分分类与回归的决策树
- 使用 DataFrame 元数据区分连续和分类特征
树集成方法
DataFrame API 支持两种主要的树集成算法:随机森林(Random Forests)和梯度提升树(GBTs),两者都使用决策树作为基础模型。
与原始 MLlib 集成 API 的主要区别:
- 支持 DataFrames 和 ML Pipelines
- 区分分类与回归
- 使用 DataFrame 元数据区分连续和分类特征
- 随机森林更多功能:特征重要性估计、分类中每个类别的预测概率
随机森林输入输出列
输入列
| 参数名 | 类型 | 默认值 | 说明 |
|---|---|---|---|
labelCol | Double | "label" | 要预测的标签 |
featuresCol | Vector | "features" | 特征向量 |
输出列(预测)
| 参数名 | 类型 | 默认值 | 说明 | 备注 |
|---|---|---|---|---|
predictionCol | Double | "prediction" | 预测标签 | |
rawPredictionCol | Vector | "rawPrediction" | 长度等于类别数的向量 | 仅分类 |
probabilityCol | Vector | "probability" | 归一化为多项分布的 rawPrediction | 仅分类 |
GBT 输入输出列
注意:
GBTClassifier当前仅支持二元标签。
| 参数名 | 类型 | 默认值 | 说明 | 备注 |
|---|---|---|---|---|
labelCol | Double | "label" | 要预测的标签 | |
featuresCol | Vector | "features" | 特征向量 | |
predictionCol | Double | "prediction" | 预测标签 |
未来
GBTClassifier也将输出rawPrediction和probability列。
聚类
K-means
K-means 是最常用的聚类算法之一,将数据点聚类到预定义数量的簇中。MLlib 实现包括 k-means++ 方法的并行化变体(称为 kmeans||)。
输入输出列
| 参数名 | 类型 | 默认值 | 说明 |
|---|---|---|---|
featuresCol | Vector | "features" | 特征向量 |
predictionCol | Int | "prediction" | 预测的簇中心 |
import org.apache.spark.ml.clustering.KMeans
import org.apache.spark.ml.evaluation.ClusteringEvaluator
// Loads data.
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains a k-means model.
val kmeans = new KMeans().setK(2).setSeed(1L)
val model = kmeans.fit(dataset)
// Make predictions
val predictions = model.transform(dataset)
// Evaluate clustering by computing Silhouette score
val evaluator = new ClusteringEvaluator()
val silhouette = evaluator.evaluate(predictions)
println(s"Silhouette with squared euclidean distance = $silhouette")
// Shows the result.
println("Cluster Centers: ")
model.clusterCenters.foreach(println)
参数说明
| 参数/方法 | 说明 |
|---|---|
setK(value: Int) | 设置要创建的簇数量 K |
setSeed(value: Long) | 设置随机种子 |
setMaxIter(value: Int) | 设置最大迭代次数 |
setTol(value: Double) | 设置收敛阈值,簇中心移动小于阈值时停止 |
setInitMode(value: String) | 设置初始化模式:random 或 `k-means |
setFeaturesCol(value: String) | 设置特征列名称 |
setPredictionCol(value: String) | 设置预测结果列名称 |
setDistanceMeasure(value: String) | 设置距离度量:euclidean(欧氏距离,默认)或 cosine(余弦相似度) |
fit(data) | 训练 KMeans 模型 |
computeCost(data) | 计算聚类成本(每个点到最近中心距离的平方和) |
transform(data) | 对数据集聚类,返回每个点的簇标签 |
summary | 获取模型摘要(簇中心、迭代次数等) |
LDA(潜在狄利克雷分配)
LDA 作为同时支持 EMLDAOptimizer 和 OnlineLDAOptimizer 的 Estimator 实现,生成 LDAModel。EMLDAOptimizer 生成的 LDAModel 可转换为 DistributedLDAModel。
import org.apache.spark.ml.clustering.LDA
// Loads data.
val dataset = spark.read.format("libsvm")
.load("data/mllib/sample_lda_libsvm_data.txt")
// Trains a LDA model.
val lda = new LDA().setK(10).setMaxIter(10)
val model = lda.fit(dataset)
val ll = model.logLikelihood(dataset)
val lp = model.logPerplexity(dataset)
println(s"The lower bound on the log likelihood of the entire corpus: $ll")
println(s"The upper bound on perplexity: $lp")
// Describe topics.
val topics = model.describeTopics(3)
println("The topics described by their top-weighted terms:")
topics.show(false)
// Shows the result.
val transformed = model.transform(dataset)
transformed.show(false)
Bisecting K-means(二分 K-means)
二分 K-means 是一种层次聚类方法,使用自顶向下的方法:所有观测从一个簇开始,随着层次向下一级一级地递归分裂。通常比常规 K-means 快得多,但一般会产生不同的聚类结果。
import org.apache.spark.ml.clustering.BisectingKMeans
// Loads data.
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains a bisecting k-means model.
val bkm = new BisectingKMeans().setK(2).setSeed(1)
val model = bkm.fit(dataset)
// Evaluate clustering.
val cost = model.computeCost(dataset)
println(s"Within Set Sum of Squared Errors = $cost")
// Shows the result.
println("Cluster Centers: ")
val centers = model.clusterCenters
centers.foreach(println)
高斯混合模型(GMM)
高斯混合模型表示一个复合分布,其中点从 K 个高斯子分布之一中抽取,每个子分布有自己的概率。spark.ml 使用期望最大化(EM)算法来学习最大似然模型。
输入输出列
| 参数名 | 类型 | 默认值 | 说明 |
|---|---|---|---|
featuresCol | Vector | "features" | 特征向量 |
predictionCol | Int | "prediction" | 预测的簇中心 |
probabilityCol | Vector | "probability" | 每个簇的概率 |
import org.apache.spark.ml.clustering.GaussianMixture
// Loads data
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains Gaussian Mixture Model
val gmm = new GaussianMixture()
.setK(2)
val model = gmm.fit(dataset)
// output parameters of mixture model model
for (i <- 0 until model.getK) {
println(s"Gaussian $i:\nweight=${model.weights(i)}\n" +
s"mu=${model.gaussians(i).mean}\nsigma=\n${model.gaussians(i).cov}\n")
}
协同过滤
ALS(交替最小二乘)
协同过滤常用于推荐系统,目标是用少量潜在因子来预测用户-物品关联矩阵中的缺失值。spark.ml 使用 ALS 算法学习这些潜在因子。
参数说明
| 参数 | 默认值 | 说明 |
|---|---|---|
numBlocks | 10 | 用户和物品分区并行计算的块数 |
rank | 10 | 模型中潜在因子的数量 |
maxIter | 10 | 最大迭代次数 |
regParam | 1.0 | ALS 中正则化参数 |
implicitPrefs | false | 是否使用隐式反馈变体 |
alpha | 1.0 | 隐式反馈变体的参数,控制偏好观察的基线置信度 |
nonnegative | false | 是否对最小二乘使用非负约束 |
注意:基于 DataFrame 的 ALS API 目前仅支持整数类型的 user 和 item ID。其他数值类型也支持但 ID 必须在整数范围内。
显式反馈 vs 隐式反馈
- 显式反馈:用户对物品的明确偏好(如评分)
- 隐式反馈:间接信号(浏览、点击、购买、点赞、分享等),数值代表用户行为的强度
对于隐式反馈,spark.ml 采用来自 Collaborative Filtering for Implicit Feedback Datasets 的方法,将数据视为代表用户行为强度的数值,模型尝试发现潜在因子来预测用户对物品的期望偏好。
正则化参数缩放
ALS 使用 “ALS-WR” 方法缩放正则化参数,使其更少依赖数据集规模,从采样子集学到的最佳参数可以应用于完整数据集。
冷启动策略
当预测时遇到训练中未出现的用户或物品时,有两种策略:
| 策略 | 说明 |
|---|---|
nan(默认) | 分配 NaN 预测值,便于系统识别新用户/物品并回退处理 |
drop | 删除包含 NaN 值的行(用于交叉验证,避免评估指标为 NaN) |
import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.recommendation.ALS
case class Rating(userId: Int, movieId: Int, rating: Float, timestamp: Long)
def parseRating(str: String): Rating = {
val fields = str.split("::")
assert(fields.size == 4)
Rating(fields(0).toInt, fields(1).toInt, fields(2).toFloat, fields(3).toLong)
}
val ratings = spark.read.textFile("data/mllib/als/sample_movielens_ratings.txt")
.map(parseRating)
.toDF()
val Array(training, test) = ratings.randomSplit(Array(0.8, 0.2))
// Build the recommendation model using ALS on the training data
val als = new ALS()
.setMaxIter(5)
.setRegParam(0.01)
.setUserCol("userId")
.setItemCol("movieId")
.setRatingCol("rating")
val model = als.fit(training)
// Evaluate the model by computing the RMSE on the test data
// Note we set cold start strategy to 'drop' to ensure we don't get NaN evaluation metrics
model.setColdStartStrategy("drop")
val predictions = model.transform(test)
val evaluator = new RegressionEvaluator()
.setMetricName("rmse")
.setLabelCol("rating")
.setPredictionCol("prediction")
val rmse = evaluator.evaluate(predictions)
println(s"Root-mean-square error = $rmse")
// Generate top 10 movie recommendations for each user
val userRecs = model.recommendForAllUsers(10)
// Generate top 10 user recommendations for each movie
val movieRecs = model.recommendForAllItems(10)
// Generate top 10 movie recommendations for a specified set of users
val users = ratings.select(als.getUserCol).distinct().limit(3)
val userSubsetRecs = model.recommendForUserSubset(users, 10)
// Generate top 10 user recommendations for a specified set of movies
val movies = ratings.select(als.getItemCol).distinct().limit(3)
val movieSubSetRecs = model.recommendForItemSubset(movies, 10)
隐式反馈使用
val als = new ALS()
.setMaxIter(5)
.setRegParam(0.01)
.setImplicitPrefs(true)
.setUserCol("userId")
.setItemCol("movieId")
.setRatingCol("rating")
频繁模式挖掘
FP-Growth
FP-Growth 算法(“FP” = Frequent Pattern)用于挖掘频繁项集。与 Apriori 类算法不同,FP-Growth 使用后缀树(FP-tree)结构编码事务,无需显式生成候选集。
spark.ml 实现了并行版本 PFP(Parallel FP-growth),基于事务后缀分布式构建 FP-tree。
参数说明
| 参数 | 说明 |
|---|---|
minSupport | 项集被视为频繁的最小支持度 |
minConfidence | 生成关联规则的最小置信度 |
numPartitions | 分布式工作的分区数(默认使用输入数据集的分区数) |
FPGrowthModel 提供的信息
| 属性/方法 | 说明 |
|---|---|
freqItemsets | 频繁项集,格式 DataFrame("items"[Array], "freq"[Long]) |
associationRules | 置信度高于 minConfidence 的关联规则,格式 DataFrame("antecedent"[Array], "consequent"[Array], "confidence"[Double]) |
transform | 对比输入项与关联规则的前件,汇总后件作为预测结果 |
import org.apache.spark.ml.fpm.FPGrowth
val dataset = spark.createDataset(Seq(
"1 2 5",
"1 2 3 5",
"1 2"
)).map(t => t.split(" ")).toDF("items")
val fpgrowth = new FPGrowth().setItemsCol("items").setMinSupport(0.5).setMinConfidence(0.6)
val model = fpgrowth.fit(dataset)
// Display frequent itemsets.
model.freqItemsets.show()
// Display generated association rules.
model.associationRules.show()
// transform examines the input items against all the association rules and summarize the
// consequents as prediction
model.transform(dataset).show()
PrefixSpan
PrefixSpan 是一种序列模式挖掘算法,用于发现频繁序列模式。
参数说明
| 参数 | 说明 |
|---|---|
minSupport | 序列模式被视为频繁的最小支持度 |
maxPatternLength | 频繁序列模式的最大长度,超过此长度的模式不会出现在结果中 |
maxLocalProjDBSize | 前缀投影数据库在开始本地迭代处理前允许的最大项数(需根据 executor 大小调优) |
sequenceCol | 数据集中序列列的名称(默认 "sequence"),该列包含 null 值的行被忽略 |
import org.apache.spark.ml.fpm.PrefixSpan
val smallTestData = Seq(
Seq(Seq(1, 2), Seq(3)),
Seq(Seq(1), Seq(3, 2), Seq(1, 2)),
Seq(Seq(1, 2), Seq(5)),
Seq(Seq(6)))
val df = smallTestData.toDF("sequence")
val result = new PrefixSpan()
.setMinSupport(0.5)
.setMaxPatternLength(5)
.setMaxLocalProjDBSize(32000000)
.findFrequentSequentialPatterns(df)
.show()
ML Tuning:模型选择与超参数调优
核心概念:Evaluator、Estimator、Pipeline 的区别
| 组件 | 说明 |
|---|---|
| Evaluator(模型评估器) | 用于评价模型的类,包括 BinaryClassificationEvaluator、MulticlassClassificationEvaluator、RegressionEvaluator、ClusteringEvaluator 等,具有 setMetricName 方法设置评估指标(如 AUC、RMSE 等) |
| Estimator(模型估计器) | 用于训练机器学习模型的抽象类,通常不具有 setMetricName 方法。例如 LogisticRegression、RandomForestClassifier、KMeans 等 |
| Pipeline(工作流) | 包含一系列 Estimator 和 Transformer 的工作流。Pipeline 本身不是评估器,因此不具有 setMetricName 方法,但可以在 Pipeline 内部的具体评估器上使用 setMetricName 方法 |
Evaluator 模型评价
BinaryClassificationEvaluator
用于评估二元分类模型性能(如逻辑回归模型),支持计算 AUC、PR AUC、准确度、F1 分数等指标。
import org.apache.spark.ml.evaluation.BinaryClassificationEvaluator
// 创建 BinaryClassificationEvaluator 实例
val evaluator = new BinaryClassificationEvaluator()
.setLabelCol("label") // 设置标签列的名称
.setRawPredictionCol("rawPrediction") // 设置原始预测值列的名称
.setMetricName("areaUnderROC") // 设置要计算的指标名称
// 使用评估器来评估模型
val auc = evaluator.evaluate(predictions) // predictions 是包含模型预测结果的 DataFrame
println(s"AUC: $auc")
常见方法
| 方法 | 说明 |
|---|---|
setLabelCol(labelCol: String) | 设置标签列的名称 |
setRawPredictionCol(rawPredictionCol: String) | 设置原始预测值列的名称 |
setMetricName(metricName: String) | 设置要计算的评估指标名称 |
evaluate(dataset: DataFrame) | 计算指定指标在给定数据集上的性能,返回 Double 类型 |
isLargerBetter | 返回布尔值,指示所选评估指标是否越大越好 |
MulticlassClassificationEvaluator
用于评估多类别分类模型性能,支持准确度、F1 分数、加权准确度等指标。
import org.apache.spark.ml.evaluation.MulticlassClassificationEvaluator
// 创建评估器实例
val evaluator = new MulticlassClassificationEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("accuracy")
// 使用评估器计算准确度
val accuracy = evaluator.evaluate(predictions)
println(s"Accuracy: $accuracy")
常见方法
| 方法 | 说明 |
|---|---|
setLabelCol(labelCol: String) | 设置标签列的名称 |
setPredictionCol(predictionCol: String) | 设置预测列的名称 |
setMetricName(metricName: String) | 设置要计算的评估指标名称 |
evaluate(dataset: DataFrame) | 计算指定指标在给定数据集上的性能 |
isLargerBetter | 返回布尔值,指示所选评估指标是否越大越好 |
RegressionEvaluator
用于评估回归模型性能,支持 MSE、RMSE、MAE、R² 等指标。
import org.apache.spark.ml.evaluation.RegressionEvaluator
// 创建评估器实例
val evaluator = new RegressionEvaluator()
.setLabelCol("label")
.setPredictionCol("prediction")
.setMetricName("rmse")
// 使用评估器计算均方根误差
val rmse = evaluator.evaluate(predictions)
println(s"RMSE: $rmse")
常见方法
| 方法 | 说明 |
|---|---|
setLabelCol(labelCol: String) | 设置标签列的名称 |
setPredictionCol(predictionCol: String) | 设置预测列的名称 |
setMetricName(metricName: String) | 设置要计算的评估指标名称 |
evaluate(dataset: DataFrame) | 计算指定指标在给定数据集上的性能 |
isLargerBetter | 返回布尔值,指示所选评估指标是否越大越好 |
ClusteringEvaluator
用于评估聚类模型性能,度量数据点分配到簇的质量。
import org.apache.spark.ml.evaluation.ClusteringEvaluator
// 创建 ClusteringEvaluator 实例
val evaluator = new ClusteringEvaluator()
.setPredictionCol("prediction")
.setFeaturesCol("features")
// 使用评估器计算 Silhouette 分数
val silhouette = evaluator.evaluate(predictions)
println(s"Silhouette: $silhouette")
常见方法
| 方法 | 说明 |
|---|---|
setPredictionCol(predictionCol: String) | 设置包含聚类预测的列的名称 |
setFeaturesCol(featuresCol: String) | 设置特征列的名称 |
evaluate(dataset: DataFrame) | 计算聚类模型在给定数据集上的性能 |
setMetricName(metricName: String) | 设置要计算的性能指标名称(支持 "silhouette") |
isLargerBetter | 返回布尔值,Silhouette 分数通常越高越好 |
setMetricName 支持的参数汇总
| 评估器 | 支持的指标 |
|---|---|
RegressionEvaluator | "rmse"(均方根误差)、"mse"(均方误差)、"r2"(R² 决定系数)、"mae"(平均绝对误差) |
BinaryClassificationEvaluator | "areaUnderROC" 或 "auc"(ROC 曲线下面积)、"areaUnderPR" 或 "auprc"(PR 曲线下面积) |
MulticlassClassificationEvaluator | "f1"(宏平均或微平均 F1 分数)、"weightedPrecision"(加权精度)、"weightedRecall"(加权召回率)、"accuracy"(准确率) |
Cross-Validation(交叉验证)
CrossValidator 通常与 ParamGridBuilder 结合使用,用于 N 折交叉验证,以选择最优的超参数组合。将数据集分成 K 个子集,每次使用其中一个子集作为验证集,其余 K-1 个子集用于训练,重复 K 次后取平均性能。
注意:
CrossValidator代价可能非常高昂,因为它需要暴力穷举指定的所有超参数组合。
import org.apache.spark.ml.Pipeline
import org.apache.spark.ml.classification.LogisticRegression
import org.apache.spark.ml.evaluation.BinaryClassificationEvaluator
import org.apache.spark.ml.feature.{HashingTF, Tokenizer}
import org.apache.spark.ml.linalg.Vector
import org.apache.spark.ml.tuning.{CrossValidator, ParamGridBuilder}
import org.apache.spark.sql.Row
// Prepare training data from a list of (id, text, label) tuples.
val training = spark.createDataFrame(Seq(
(0L, "a b c d e spark", 1.0),
(1L, "b d", 0.0),
(2L, "spark f g h", 1.0),
(3L, "hadoop mapreduce", 0.0),
(4L, "b spark who", 1.0),
(5L, "g d a y", 0.0),
(6L, "spark fly", 1.0),
(7L, "was mapreduce", 0.0),
(8L, "e spark program", 1.0),
(9L, "a e c l", 0.0),
(10L, "spark compile", 1.0),
(11L, "hadoop software", 0.0)
)).toDF("id", "text", "label")
// Configure an ML pipeline, which consists of three stages: tokenizer, hashingTF, and lr.
val tokenizer = new Tokenizer()
.setInputCol("text")
.setOutputCol("words")
val hashingTF = new HashingTF()
.setInputCol(tokenizer.getOutputCol)
.setOutputCol("features")
val lr = new LogisticRegression()
.setMaxIter(10)
val pipeline = new Pipeline()
.setStages(Array(tokenizer, hashingTF, lr))
// We use a ParamGridBuilder to construct a grid of parameters to search over.
// With 3 values for hashingTF.numFeatures and 2 values for lr.regParam,
// this grid will have 3 x 2 = 6 parameter settings for CrossValidator to choose from.
val paramGrid = new ParamGridBuilder()
.addGrid(hashingTF.numFeatures, Array(10, 100, 1000))
.addGrid(lr.regParam, Array(0.1, 0.01))
.build()
// We now treat the Pipeline as an Estimator, wrapping it in a CrossValidator instance.
// A CrossValidator requires an Estimator, a set of Estimator ParamMaps, and an Evaluator.
val cv = new CrossValidator()
.setEstimator(pipeline)
.setEvaluator(new BinaryClassificationEvaluator)
.setEstimatorParamMaps(paramGrid)
.setNumFolds(2) // Use 3+ in practice
.setParallelism(2) // Evaluate up to 2 parameter settings in parallel
// Run cross-validation, and choose the best set of parameters.
val cvModel = cv.fit(training)
// Prepare test documents, which are unlabeled (id, text) tuples.
val test = spark.createDataFrame(Seq(
(4L, "spark i j k"),
(5L, "l m n"),
(6L, "mapreduce spark"),
(7L, "apache hadoop")
)).toDF("id", "text")
// Make predictions on test documents. cvModel uses the best model found (lrModel).
cvModel.transform(test)
.select("id", "text", "probability", "prediction")
.collect()
.foreach { case Row(id: Long, text: String, prob: Vector, prediction: Double) =>
println(s"($id, $text) --> prob=$prob, prediction=$prediction")
}
参数说明
| 参数/方法 | 说明 |
|---|---|
setEstimator | 设置要评估的机器学习模型(如 RandomForest、LogisticRegression、Pipeline 等) |
setEstimatorParamMaps | 设置要尝试的不同超参数组合的参数网格,通过 ParamGridBuilder 构建 |
setEvaluator | 设置用于评估模型性能的评估器 |
setNumFolds | 设置 K 折交叉验证中的 K 值(默认 3) |
setParallelism | 设置交叉验证的并行度 |
fit(dataset) | 执行 K 折交叉验证,返回包含每个超参数组合性能度量值的 CrossValidatorModel |
ParamGridBuilder 方法
| 方法 | 说明 |
|---|---|
addGrid(param: Param[_], values: Array[_]) | 添加要调整的参数及其可能的取值范围 |
build() | 构建参数网格,返回 Array[ParamMap],每个 ParamMap 表示一个参数组合 |
Train-Validation Split(训练验证分割)
除了 CrossValidator 之外,Spark 还提供 TrainValidationSplit 用于超参数调优。TrainValidationSplit 只评估每个参数组合一次(而非 K 次),因此更经济,但当训练数据集不够大时不会产生可靠的结果。
使用 trainRatio 参数将数据集拆分为训练集和验证集。例如 trainRatio=0.75 时,75% 数据用于训练,25% 用于验证。与 CrossValidator 一样,TrainValidationSplit 最终使用最好的 ParamMap 和整个数据集来拟合估算器。
import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.regression.LinearRegression
import org.apache.spark.ml.tuning.{ParamGridBuilder, TrainValidationSplit}
// Prepare training and test data.
val data = spark.read.format("libsvm").load("data/mllib/sample_linear_regression_data.txt")
val Array(training, test) = data.randomSplit(Array(0.9, 0.1), seed = 12345)
val lr = new LinearRegression()
.setMaxIter(10)
// We use a ParamGridBuilder to construct a grid of parameters to search over.
val paramGrid = new ParamGridBuilder()
.addGrid(lr.regParam, Array(0.1, 0.01))
.addGrid(lr.fitIntercept)
.addGrid(lr.elasticNetParam, Array(0.0, 0.5, 1.0))
.build()
// A TrainValidationSplit requires an Estimator, a set of Estimator ParamMaps, and an Evaluator.
val trainValidationSplit = new TrainValidationSplit()
.setEstimator(lr)
.setEvaluator(new RegressionEvaluator)
.setEstimatorParamMaps(paramGrid)
.setTrainRatio(0.8) // 80% of the data will be used for training
.setParallelism(2) // Evaluate up to 2 parameter settings in parallel
// Run train validation split, and choose the best set of parameters.
val model = trainValidationSplit.fit(training)
// Make predictions on test data. model is the model with combination of parameters
// that performed best.
model.transform(test)
.select("features", "label", "prediction")
.show()
参数说明
| 参数/方法 | 说明 |
|---|---|
setEstimator | 设置要评估的机器学习模型 |
setEstimatorParamMaps | 设置要尝试的不同超参数组合的参数网格 |
setEvaluator | 设置用于评估模型性能的评估器 |
setTrainRatio | 设置用于训练的数据比例(0 到 1 之间,默认 0.75) |
setParallelism | 设置并行度 |
fit(dataset) | 执行训练验证划分,返回 TrainValidationSplitModel,可通过 bestModel 查看最优模型 |
CrossValidator vs TrainValidationSplit
| CrossValidator | TrainValidationSplit | |
|---|---|---|
| 评估次数 | K 次(每个参数组合) | 1 次(每个参数组合) |
| 可靠性 | 高 | 依赖数据集大小 |
| 计算代价 | 高(暴力穷举) | 较低 |
| 适用场景 | 数据集不是特别大 | 大数据集,快速迭代 |
Pipeline Stage 模型自定义
Transformer 自定义
Transformer 对数据做转换操作,不涉及模型训练,一般为中间结果,不需要保存到本地。效果为在原有 DataFrame 上增加一列。
创建步骤
- 继承
Transformer抽象类 - [可选] 实现
setInputCol和setOutputCol方法 - 实现
transformSchema方法 - 实现
transform方法 - [可选] 实现可读写
- 通过 Pipeline 反射调用
步骤 1:[可选] 实现 setInputCol 和 setOutputCol 方法
初始化 inputCol 和 outputCol 属性(类型为 Param[String]),setter 方法设置属性后需返回 this(建造者模式)。需将 val 改为 var 以便 setter 修改。
import org.apache.spark.ml.param.{Param, Params, ParamMap}
import org.apache.spark.ml.util.Identifiable
import org.apache.spark.ml.Transformer
class Mytransformer(override val uid: String) extends Transformer with Identifiable {
final var inputCol = new Param[String](this, "inputCol", "默认输入列")
final var outputCol = new Param[String](this, "outputCol", "默认输出列")
def setInputCol(value: String): this.type = {
this.inputCol = new Param[String](this, value, "指定输入列")
this
}
def setOutputCol(value: String): this.type = {
this.outputCol = new Param[String](this, value, "指定输出列")
this
}
// 辅助构造函数,支持不传值初始化
def this() = this(Identifiable.randomUID("Mytransformer"))
// 继承 copy 函数,使用默认实现即可
// def copy(extra: ParamMap): Mytransformer = defaultCopy(extra)
}
步骤 2:重写 transformSchema 方法
transformSchema 用于改变 DataFrame 的 schema,对输入列进行检查,并增加输出的列到 schema。
override def transformSchema(schema: StructType): StructType = {
// 可以加入字段校验逻辑
val idx = schema.fieldIndex(this.inputCol.name)
val field = schema.fields(idx)
if (field.dataType != DoubleType) {
throw new Exception(s"字段${field.name}输入类型${field.dataType}与预期类型DoubleType不符")
}
// 转换原 DataFrame 的 schema
schema.add(StructField(outputCol.name, DoubleType, false))
}
步骤 3:实现 transform 方法
override def transform(dataset: Dataset[_]): DataFrame = {
// 数据处理逻辑,DataFrame 常见操作
dataset.withColumn(outputCol.name, round(col(inputCol.name)))
}
步骤 4:[可选] 实现可读写
让类继承 DefaultParamsWritable:
class Mytransformer(override val uid: String) extends Transformer with DefaultParamsWritable
实现 Transformer 的伴生对象,继承 DefaultParamsReadable,重写 load 方法:
object Mytransformer extends DefaultParamsReadable[Mytransformer] {
override def load(path: String): Mytransformer = super.load(path)
}
步骤 5:通过 Pipeline 使用自定义 Transformer
val dataset = spark.createDataFrame(Seq(
("mike", 166.0), ("tom", 175.0), ("wade", 163.0)
)).toDF("name", "height")
val mytransformer = new Mytransformer
mytransformer.setInputCol("height").setOutputCol("h170")
val pipeline = new Pipeline().setStages(Array(mytransformer)).fit(dataset)
val r = pipeline.transform(dataset)
r.show()
Estimator 自定义
Estimator 一般做预测(需要数据的历史信息),涉及模型训练,需要保存结果到本地。自定义 Estimator 类需要两部分内容:自定义 Model(继承 Model)和自定义 Estimator(继承 Estimator)。
创建步骤
- 构造通用参数 Trait
- 创建自定义
EstimatorModel类(继承Model、混入参数 Trait,重写transformSchema和transform) - 创建自定义
Estimator类(继承Estimator、混入参数 Trait,重写transformSchema和fit) - 实现可读写
- 通过 Pipeline 调用
步骤 1:创建通用参数 Trait
自定义 EstimatorModel 类和自定义 Estimator 类拥有相同的模型参数,抽象出用于表示参数的 Trait 可防止代码重复。
trait MyEstimatorParams extends Params {
final var inputCol = new Param[String](this, "inputCol", "默认输入列")
final var outputCol = new Param[String](this, "outputCol", "默认输出列")
def setInputCol(value: String) = {
this.inputCol = new Param[String](this, value, "指定输入列")
this
}
def setOutputCol(value: String) = {
this.outputCol = new Param[String](this, value, "指定输出列")
this
}
}
步骤 2:创建自定义 EstimatorModel 类
class MyEstimatorModel(override val uid: String, val dataset: Dataset[_])
extends Model[MyEstimatorModel] with MyEstimatorParams {
override def copy(extra: ParamMap): MyEstimatorModel = {
val copied = new MyEstimatorModel(uid, dataset)
// parent 属性继承自抽象类 Model
copyValues(copied, extra).setParent(parent)
}
}
重写 transformSchema
override def transformSchema(schema: StructType): StructType = {
val idx = schema.fieldIndex(this.inputCol.name)
val field = schema.fields(idx)
if (field.dataType != DoubleType) {
throw new Exception(s"输入列${this.inputCol.name}类型${field.dataType}与预期类型DoubleType不符")
}
schema.add(StructField(this.outputCol.name, DoubleType, false))
}
实现 transform(结合 UDF 函数)
override def transform(dataset: Dataset[_]): DataFrame = {
val average = dataset.groupBy().agg(avg(col(this.inputCol.name))).collect().apply(0).getDouble(0)
val standard = dataset.groupBy().agg(stddev_pop(col(this.inputCol.name))).collect().apply(0).getDouble(0)
// 剔除异常值:与均值的差值,超出 3 倍标准差
val function = udf((value: Double) => if (value - average > 3 * standard) 0 else value)
dataset.withColumn(this.outputCol.name, function(col(this.inputCol.name)))
}
步骤 3:创建自定义 Estimator 类
class MyEstimator(override val uid: String)
extends Estimator[MyEstimatorModel] with MyEstimatorParams {
def this() = this(Identifiable.randomUID("myEstimatorParams"))
override def copy(extra: ParamMap): MyEstimator = defaultCopy(extra)
}
重写 transformSchema
override def transformSchema(schema: StructType): StructType = {
val idx = schema.fieldIndex(this.inputCol.name)
val field = schema.fields(idx)
if (field.dataType != DoubleType) {
throw new Exception(s"输入列${this.inputCol.name}类型${field.dataType}与预期类型DoubleType不符")
}
schema.add(StructField(this.outputCol.name, DoubleType, false))
}
重写 fit 方法
fit 方法仅向 MyEstimatorModel 传参,指定输入列、输出列。
override def fit(dataset: Dataset[_]): MyEstimatorModel = {
val c = new MyEstimatorModel(uid, dataset)
c.setInputCol(this.inputCol.name).setOutputCol(this.outputCol.name)
c
}
步骤 4:实现可读写
由于 Spark 源码中 DefaultParamsWriter 和 DefaultParamsReader 的方法是私有的,自定义类无法直接通过继承获得。需要从源码 spark.ml.util.ReadWrite.scala 中提取相关内容,实现 MyDefaultParamsWriter 和 MyDefaultParamsReader 两个 Object。
主要步骤:
- EstimatorModel 伴生对象:继承
MLReadable,实现内部类MyEstimatorReader(继承MLReader,重写load方法) - EstimatorModel 类更新:继承
MLWritable,实现内部类MyEstimatorWriter,重写write方法 - Estimator 类更新:继承
DefaultParamsReadable,实现 Estimator 伴生对象
步骤 5:通过 Pipeline 调用
val dataset = spark.createDataFrame(Seq(
("mike", 166.0), ("tom", 175.0), ("wade", 163.0), ("bad", 660.0),
("james", 160.0), ("black", 166.0), ("angel", 162.0), ("Emma", 177.0),
("weekn", 174.0), ("kelly", 166.0), ("grey", 140.0)
)).toDF("name", "height")
val est = new MyEstimator
est.setInputCol("height").setOutputCol("normal")
val pipeline = new Pipeline().setStages(Array(est))
val model = pipeline.fit(dataset)
model.write.overwrite().save("model") // 如果未继承实现 MLWritable,则会报错
val res = model.transform(dataset)
res.show()
完整示例:自定义 WOE Estimator
WOE(Weight of Evidence,证据权重)是一种常用的特征工程方法,用于将分类变量转换为连续变量。下面展示一个完整的自定义 Estimator 实现,包括参数 Trait、Model 类、Estimator 类以及可读写的完整实现。
参数 Trait:WOEBase
管理输入列、输出列、标签列、阈值等参数。
import org.apache.spark.ml.util._
import org.apache.spark.ml.param._
import org.apache.spark.ml.attribute._
import org.apache.spark.ml.Estimator
import org.apache.spark.ml.Model
import org.apache.spark.sql.{ DataFrame, Dataset }
import org.apache.spark.sql.types._
import org.apache.spark.sql.functions._
import scala.collection.mutable.ArrayBuffer
import org.apache.spark.SparkException
// 特质:管理输入列、输出列、标签列、阈值等参数
trait WOEBase extends Params {
// 从 Spark 2.3 开始,可以直接继承 sharedParams 里面的相关特质
final val labelCol: Param[String] = new Param[String](this, "labelCol", "label column name")
final val inputCol: Param[String] = new Param[String](this, "inputCol", "input column name")
final val outputCol: Param[String] = new Param[String](this, "outputCol", "output column name")
final val inputCols: StringArrayParam = new StringArrayParam(this, "inputCols", "input column names")
final val outputCols: StringArrayParam = new StringArrayParam(this, "outputCols", "output column names")
// $ 是 Params 中实现的方法,实际上调用了 getOrDefault 方法
def getLabelCol() = $(labelCol)
def getInputCol() = $(inputCol)
def getOutputCol() = $(outputCol)
def getInputCols() = $(inputCols)
def getOutputCols() = $(outputCols)
final val delta: DoubleParam = new DoubleParam(this, "delta", "防止出现0值,造成除0溢出或对数无穷大,而增加的修正值")
def getDelta() = $(delta)
// -> 是 Params 中实现的方法,用于生成一个 ParamPair
setDefault(delta -> 1, labelCol -> "label") // Params 类的方法,设置参数默认值
参数说明
| 参数 | 说明 |
|---|---|
labelCol | 标签列名称 |
inputCol | 单个输入列名称 |
outputCol | 单个输出列名称 |
inputCols | 多个输入列名称数组(与 outputCols 对应) |
outputCols | 多个输出列名称数组(与 inputCols 对应) |
delta | 防止出现 0 值的修正值(默认 1) |
$(param)等价于getOrDefault(param),用于获取参数值;isSet(param)检查参数是否被set方法设置过(通过setDefault设置的默认参数不返回true)。
输入输出列处理逻辑
// 格式化输入、输出列,统一返回 (Array(inputCol), Array(outputCol))
protected def getInOutCols: (Array[String], Array[String]) = {
// require 方法是 scala.Predef 对象下的预定义方法,条件为 false 则抛出异常
require(
(isSet(inputCol) && isSet(outputCol) && !isSet(inputCols) && !isSet(outputCols)) ||
(!isSet(inputCol) && !isSet(outputCol) && isSet(inputCols) && isSet(outputCols)),
"WOE only supports setting either inputCol/outputCol or inputCols/outputCols.")
if (isSet(inputCol)) {
(Array($(inputCol)), Array($(outputCol)))
} else {
require($(inputCols).length == $(outputCols).length,
"inputCols number do not match outputCols")
($(inputCols), $(outputCols))
}
}
Schema 验证方法
protected def validateAndTransformSchemas(schema: StructType): StructType = {
val labelColName = $(labelCol)
val labelDataType = schema(labelColName).dataType
require(
labelDataType.isInstanceOf[NumericType],
s"The label column $labelColName must be numeric type, but got $labelDataType.")
val (inputColNames, outputColNames) = getInOutCols
val existingFields = schema.fields
var outputFields = existingFields
inputColNames.zip(outputColNames).foreach {
case (inputColName, outputColName) =>
require(existingFields.exists(_.name == inputColName),
s"Input column ${inputColName} not exists.")
require(existingFields.forall(_.name != outputColName),
s"Output column ${outputColName} already exists.")
val attr = NominalAttribute.defaultAttr.withName(outputColName)
outputFields :+= attr.toStructField()
}
StructType(outputFields)
}
Model 类:WOEModel
实现 transform 方法,并继承 MLWritable。
class WOEModel(override val uid: String, val woe_map_arr: Seq[Map[String, Double]])
extends Model[WOEModel]
with MLWritable with WOEBase {
def this(woe_map_arr: Seq[Map[String, Double]]) =
this(Identifiable.randomUID("WOE"), woe_map_arr)
def setLabelCol(value: String): this.type = set(labelCol, value)
def setInputCol(value: String): this.type = set(inputCol, value)
def setOutputCol(value: String): this.type = set(outputCol, value)
def setInputCols(value: Array[String]): this.type = set(inputCols, value)
def setOutputCols(value: Array[String]): this.type = set(outputCols, value)
// 模型不可以设置 delta,因为 delta 只对学习有用,对转换没用
override def copy(extra: ParamMap): WOEModel = {
val copied = new WOEModel(uid, woe_map_arr)
copyValues(copied, extra).setParent(parent)
}
import WOEModel._
override def write: WOEModelWriter = new WOEModelWriter(this)
override def transform(dataset: Dataset[_]): DataFrame = {
val (inputColNames, outputColNames) = getInOutCols
transformSchema(dataset.schema)
require(woe_map_arr.length == inputColNames.length,
s"The number of input columns is not equal to the number of WOEModel model maps")
var df: DataFrame = dataset.toDF()
woe_map_arr.zipWithIndex.map {
case (woe_map, idx) =>
val inputColName = inputColNames(idx)
val outputColName = outputColNames(idx)
val woer = udf { (feature: String) =>
woe_map.get(feature) match {
case Some(n: Double) => n
case None =>
throw new SparkException(
s"Input column_${inputColName}'s value ${feature} " +
"does not exist in the WOEModel model map.")
}
}
df = df.withColumn(outputColName, woer(dataset(inputColName).cast(StringType)))
}
df
}
override def transformSchema(schema: StructType): StructType = {
validateAndTransformSchemas(schema)
}
}
Model 伴生对象:实现读取功能
object WOEModel extends MLReadable[WOEModel] {
import org.apache.hadoop.fs.Path
import org.json4s.JsonDSL._
import org.json4s.jackson.JsonMethods._
import org.json4s.JsonAST._
implicit val format = org.json4s.DefaultFormats
// 自定义 Writer
private[WOEModel] class WOEModelWriter(instance: WOEModel) extends MLWriter {
private case class Data(woe_map_arr: Seq[Map[String, Double]])
override protected def saveImpl(path: String): Unit = {
val metadataPath = new Path(path, "metadata").toString
val params = instance.extractParamMap().toSeq.asInstanceOf[Seq[ParamPair[Any]]]
val jsonParams = render(params.map {
case ParamPair(p, v) => p.name -> parse(p.jsonEncode(v))
}.toList)
val basicMetadata = ("class" -> instance.getClass.getName) ~
("timestamp" -> System.currentTimeMillis()) ~
("sparkVersion" -> sc.version) ~
("uid" -> instance.uid) ~
("paramMap" -> jsonParams)
val metadataJson = compact(render(basicMetadata))
sc.parallelize(Seq(metadataJson), 1).saveAsTextFile(metadataPath)
val data = Data(instance.woe_map_arr)
val dataPath = new Path(path, "data").toString
sparkSession.createDataFrame(Seq(data)).repartition(1).write.parquet(dataPath)
}
}
// 自定义 Reader
private class WOEModelReader extends MLReader[WOEModel] {
private val className = classOf[WOEModel].getName
override def load(path: String): WOEModel = {
val metadataPath = new Path(path, "metadata").toString
val s = sc.textFile(metadataPath, 1).first()
val metadata = parse(s)
val clz = (metadata \ "class").extract[String]
val uid = (metadata \ "uid").extract[String]
require(className == clz,
s"Error loading metadata: Expected class name $className but found $clz")
val dataPath = new Path(path, "data").toString
val data = sparkSession.read.parquet(dataPath).select("woe_map_arr").head()
val woe_map_arr = data.getAs[Seq[Map[String, Double]]](0)
val instance = new WOEModel(uid, woe_map_arr)
val params = metadata \ "paramMap"
params match {
case JObject(pairs) =>
pairs.foreach {
case (paramName, jsonValue) =>
val param = instance.getParam(paramName)
val value = param.jsonDecode(compact(render(jsonValue)))
instance.set(param, value)
}
case _ =>
throw new IllegalArgumentException(s"Cannot recognize JSON metadata: ${s}.")
}
instance
}
}
override def read: MLReader[WOEModel] = new WOEModelReader
override def load(path: String): WOEModel = super.load(path)
}
Estimator 类:WOE
实现模型训练方法 fit 和写功能。
架构说明:由于不支持多重继承,因此无法同时继承读(
DefaultParamsReadable)和写(DefaultParamsWritable)的默认实例类。解决方案:让类继承并实现写的父类(DefaultParamsWritable),让伴生对象继承读的父类(DefaultParamsReadable)。
class WOE(override val uid: String)
extends Estimator[WOEModel]
with WOEBase
with DefaultParamsWritable {
def this() = this(Identifiable.randomUID("WOE"))
// set 方法
def setLabelCol(value: String): this.type = set(labelCol, value)
def setInputCol(value: String): this.type = set(inputCol, value)
def setOutputCol(value: String): this.type = set(outputCol, value)
def setInputCols(value: Array[String]): this.type = set(inputCols, value)
def setOutputCols(value: Array[String]): this.type = set(outputCols, value)
def setDelta(value: Double): this.type = set(delta, value)
override def copy(extra: ParamMap): this.type = defaultCopy(extra)
override def fit(dataset: Dataset[_]): WOEModel = {
transformSchema(dataset.schema, true)
val delta_value = $(delta) // 防止出现 0 值,而增加的修正
val T = dataset.count
val B = dataset.where($(labelCol) + " = 1").count()
val G = T - B
val woe_map_arr = new ArrayBuffer[Map[String, Double]]()
val (inputColNames, outputColNames) = getInOutCols
inputColNames.foreach { inputColName =>
val gDs_t = dataset.groupBy(inputColName)
.agg(count($(labelCol)).as("T"), sum($(labelCol)).as("B"))
val gDs = gDs_t.withColumn("G", gDs_t("T") - gDs_t("B"))
val loger = udf { d: Double => math.log(d) }
// WOE = ln((B_i + Δ) / (B + Δ) * (G + Δ) / (G_i + Δ))
val woe_map = gDs.withColumn("woe",
loger((gDs("B") + delta_value) / (B + delta_value) *
(G + delta_value) / (gDs("G") + delta_value)))
.select(col(inputColName).cast(StringType), col("woe"))
.collect()
.map(r => (r.getString(0), r.getDouble(1)))
.toMap
woe_map_arr += woe_map
}
// copyValues:将 parent 的参数值拷贝给 model
copyValues(new WOEModel(uid, woe_map_arr.toSeq).setParent(this))
}
override def transformSchema(schema: StructType): StructType = {
validateAndTransformSchemas(schema)
}
}
Estimator 伴生对象
// save 方法调用的就是 write.save,load 方法调用的是 read.load
object WOE extends DefaultParamsReadable[WOE] {
override def load(path: String): WOE = super.load(path)
}
核心原理:Estimator 学习输出 Transformer 实际上就是传递一个数据结构。
fit方法将学到的数据结果传给 Transformer——可以直接作为构造参数传递,也可以用设置参数的形式传递。这里采用构造参数传递,并简化逻辑将单列和多列统一当作多列处理。
测试代码
val train = List(
("0", "t1"), ("0", "t2"), ("0", "t3"), ("0", "t4"), ("0", "t5"),
("1", "t6"), ("1", "t7"), ("1", "t8"), ("1", "t9"), ("1", "t10")
).toDF("label", "uid")
val modelPath = "/path/to/model"
val resultPath = "/path/to/result"
// 测试训练功能
val model = new WOE().setLabelCol("label").setInputCol("uid").setOutputCol("output").fit(train)
// 测试写入功能
model.write.overwrite().save(modelPath)
// 测试预测功能
model.transform(train).show(false)
// 测试写入结果
model.transform(train).write.csv(resultPath)
// 测试读取功能
WOEModel.load(modelPath).setLabelCol("label").setInputCol("uid")
.setOutputCol("output").transform(train).show(false)
Pipeline 集成元模型
Pipeline 可以对已训练 stage 和未训练 stage 进行组合,且不会改变已训练 stage。
from pyspark.ml.pipeline import Pipeline
from pyspark.ml.feature import StringIndexer
from pyspark.ml.feature import StandardScaler
from pyspark.ml.feature import VectorAssembler
# 创建两个 df:df0 用于训练 StringIndexer,df1 用于测试
df0 = spark.createDataFrame([
['one', 1.0], ['two', 2.0], ['three', 3.0],
['one', 2.0], ['two', 1.0], ['three', 2.0],
['one', 3.0], ['two', 3.0], ['three', 2.0]
], ['c', "i"])
df1 = spark.createDataFrame([
['one', 1.0], ['two', 2.0], ['three', 3.0], ['four', 4.0]
], ['c', "i"])
# 使用 df0 训练一个 StringIndexer,后续在 Pipeline 中作为已训练 stage
string_indexer = StringIndexer() \
.setInputCol("c").setOutputCol("feature") \
.setHandleInvalid("skip").fit(df0)
# 未训练 stage-1 和 stage-2
vector_assembler = VectorAssembler() \
.setInputCols(["feature", "i"]).setOutputCol("features")
standard = StandardScaler() \
.setInputCol("features").setOutputCol("standarded")
# 创建 Pipeline,加入已训练 stage 和未训练 stage
pipeline = Pipeline().setStages([string_indexer, vector_assembler, standard])
# 使用新数据集 df1 对 Pipeline 进行训练
# 注意:已训练 stage 不会被再次训练,无法处理 df1 中的 "four" 值
model = pipeline.fit(df1)
# 处理 df0:可以处理所有记录(已训练 stage 是根据 df0 得来的)
result0 = model.transform(df0)
# 处理 df1:"four" 的记录被丢弃(已训练 stage 不可变)
result1 = model.transform(df1)
关键结论:
- 已训练 stage 在 Pipeline 中不会被再次训练,参数维持不变
- 已训练 stage 处理未知值时可能丢弃对应行(取决于
handleInvalid设置)- 通过
pipeline.getStages()[0].extractParamMap()和model.stages[0].extractParamMap()可验证参数一致性
linalg 线性代数库
linalg 库(Linear Algebra Library)是 Spark MLlib 中的一个模块,支持 MLlib 的底层数据结构,提供向量、矩阵及其运算工具。
BLAS(Basic Linear Algebra Subprograms)
BLAS 为 private[spark],只有 spark 包内部成员才能调用,提供高性能矩阵和向量计算底层接口。
| 方法 | 说明 | 签名 |
|---|---|---|
gemm | 矩阵-矩阵乘法 | gemm(alpha, A, B, beta, C) |
gemv | 矩阵-向量乘法 | gemv(trans, alpha, A, x, beta, y) |
dot | 向量点积 | dot(x: Vector, y: Vector): Double |
scal | 向量标量乘法 | scal(a: Double, x: Vector): Unit |
axpy | 向量加法 y = a*x + y | axpy(a: Double, x: Vector, y: Vector): Unit |
Matrix(矩阵接口)
Matrix 是 DenseMatrix 和 SparseMatrix 的父类,用于访问和操作矩阵元素。
| 方法 | 说明 |
|---|---|
numRows | 获取矩阵的行数 |
numCols | 获取矩阵的列数 |
apply(i, j) | 获取矩阵中特定位置的元素,索引从 0 开始 |
toArray | 将矩阵转换为二维数组 |
transpose | 返回矩阵的转置 |
multiply(anotherMatrix) | 矩阵乘法,实际调用 BLAS.gemm |
equals(anotherMatrix) | 判断两个矩阵是否相等 |
Matrices(矩阵工厂对象)
用于创建不可变矩阵的对象工厂。
稠密矩阵创建
| 方法 | 说明 |
|---|---|
dense(numRows, numCols, values) | 创建稠密矩阵,values 按列主序排列 |
diag(vector) | 创建稠密对角线矩阵 |
eye(n) | 创建稠密单位矩阵 |
zeros(n) | 创建零矩阵 |
ones(m, n) | 创建元素全为 1 的矩阵 |
horzcat(matrices) | 横向拼接矩阵(行数一致) |
vertcat(matrices) | 纵向拼接矩阵(列数一致) |
rand(m, n, rng) | 创建 [0,1) 均匀分布的随机数矩阵 |
randn(m, n, rng) | 创建均值为 0、方差为 1 的随机数矩阵 |
稀疏矩阵创建
| 方法 | 说明 |
|---|---|
sparse(numRows, numCols, colPtr, rowIndex, values) | 创建稀疏矩阵(CSC 格式) |
speye(n) | 创建稀疏单位矩阵 |
sprand(m, n, density, rng) | 创建 [0,1) 均匀分布的稀疏随机数矩阵 |
sprandn(m, n, density, rng) | 创建均值为 0、方差为 1 的稀疏随机数矩阵 |
稀疏矩阵 CSC 格式参数说明:
colPtr:长度为numCols + 1的单调递增数组,表示每列非零元素在values中的索引区间。如[0, 2, 2, 4, 4]表示第 1 列区间[0,2),第 2 列[2,2)(无元素),第 3 列[2,4)rowIndex:每个非零元素的行索引,长度与values相同values:按列主序排列的非零元素值数组
DenseMatrix(稠密矩阵)
val data = Array(1.0, 2.0, 3.0, 4.0, 5.0, 6.0) // 按先列后行顺序遍历
val dense = new DenseMatrix(3, 2, data) // 3 行 2 列
| 方法 | 说明 |
|---|---|
numRows / numCols | 返回行数/列数 |
apply(i, j) | 获取 (i, j) 位置的元素 |
toArray | 以一维数组形式返回所有元素 |
transpose | 返回矩阵的转置 |
multiply(B) | 矩阵乘法 |
copy() | 创建矩阵副本 |
colIter(j) | 返回指定列的 DenseVector 迭代器 |
rowIter(i) | 返回指定行的 DenseVector 迭代器 |
SparseMatrix(稀疏矩阵)
val data = Array(1.0, 2.0, 3.0, 4.0, 5.0, 6.0)
val colPtr = Array(0, 3, 6)
val rowIndices = Array(0, 1, 2, 0, 1, 2)
val sparse = new SparseMatrix(3, 2, colPtr, rowIndices, data)
colPtr元素数量必须为numCols + 1,第一个元素必须为 0,最后一个元素必须为data.length;rowIndices长度必须与data一致。
方法同 DenseMatrix,不同之处在于 colIter/rowIter 返回的是 SparseVector。
Vector(向量接口)
Vector 是 DenseVector 和 SparseVector 的父类。
| 方法 | 说明 |
|---|---|
size | 返回向量大小 |
toArray | 将向量转换为数组 |
apply(i) | 获取指定索引位置的元素 |
toSparse | 将稠密向量转换为稀疏向量 |
compressed | 压缩稀疏向量中零元素 |
toDense | 将稀疏向量转换为稠密向量 |
Vectors(向量工厂对象)
| 方法 | 说明 |
|---|---|
dense(values*) | 创建稠密向量 |
sparse(size, indices, values) | 创建稀疏向量 |
zeros(size) | 创建稠密零向量 |
size | 返回向量大小 |
norm(p) | 计算向量范数(p=1 为 L1,p=2 为 L2) |
sqdist(v1, v2) | 计算向量欧几里德距离平方 |
Vector UDF 操作示例
from random import random
from functools import partial
from pyspark.sql.functions import *
from pyspark.sql.types import *
from pyspark.ml.linalg import Vector, Vectors, VectorUDT
# 替换向量中指定位置的元素为随机数
def replace_dimension_with_random(dimension_to_replace, vector):
values = vector.toArray()
values[dimension_to_replace] = random()
return Vectors.dense(*values)
replace_dimension_with_random_udf = lambda dim: \
udf(lambda v: partial(replace_dimension_with_random, dim)(v), VectorUDT())
# 获取向量中指定位置的元素
element_at = lambda dimension: udf(lambda vector: str(vector[dimension]), StringType())
# 合并多个向量列,返回稀疏向量
def vector_merge_sparse(*vectors):
elements = []
for vector in vectors:
elements.extend(vector)
return Vectors.sparse(len(elements),
[(i, v) for i, v in enumerate(elements) if v != 0])
vector_merge_sparse_udf = udf(vector_merge_sparse, VectorUDT())
# 合并多个向量列,返回稠密向量
def vector_merge_dense(*vectors):
elements = []
for vector in vectors:
elements.extend(vector)
return Vectors.dense(*elements)
vector_merge_dense_udf = udf(vector_merge_dense, VectorUDT())