Article

数据计算 Spark MLlib

更新于:2026-07-13

概述

MLlib 是 Spark 的机器学习(ML)库。其目标是使实用的机器学习可扩展且容易。在较高级别,它提供了以下工具:

  • ML 算法:常见的学习算法,例如分类、回归、聚类和协作过滤
  • 特征化:特征提取、变换、降维和选择
  • 管道:用于构建、评估和调整 ML 管道的工具
  • 持久性:保存和加载算法、模型和管道
  • 实用程序:线性代数、统计信息、数据处理等

依赖库

MLlib 使用线性代数程序包 Breeze,该程序依赖于 netlib-java 进行优化的数值处理。

相关性

计算两个系列数据之间的相关性是统计中的常见操作。spark.ml 提供了很多系列中的灵活性,计算两两相关性。目前支持的相关方法是 PearsonSpearman 的相关。

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:将多个 TransformerEstimator 链接在一起以指定 ML 工作流程。
  • Parameter:所有 TransformerEstimator 共享一个用于指定参数的通用 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。例如,学习算法 LogisticRegressionEstimator,调用 fit() 训练出 LogisticRegressionModel,其为 Model,也是 Transformer

管道组件的属性

Transformer.transform()Estimator.fit() 都是无状态的。将来可能通过替代概念支持有状态算法。

TransformerEstimator 的每个实例都有一个唯一的 ID,在指定参数时很有用。

Pipeline

在机器学习中,通常需要按顺序运行一系列算法来处理数据并从中学习。例如,一个简单的文本文档处理工作流程可能包括几个阶段:

  1. 将每个文档的文本拆分为单词
  2. 将每个文档的单词转换为数字特征向量
  3. 使用特征向量和标签学习预测模型

MLlib 将这样的工作流程表示为 Pipeline,其中包含要按特定顺序运行的一系列 PipelineStageTransformerEstimator)。

工作原理

Pipeline 被指定为一个阶段序列,每个阶段都是一个 Transformer 或一个 Estimator。这些阶段按顺序运行,输入 DataFrame 在通过每个阶段时都会进行转换:

  • 对于 Transformer 阶段,transform() 方法作用于 DataFrame
  • 对于 Estimator 阶段,fit() 作用于 DataFrame,生成 transform() 方法

在训练时,Pipeline.fit() 在原始 DataFrame 上调用,经过各个阶段:

  1. Tokenizer.transform() 将原始文本文档拆分为单词,添加带有单词的新列
  2. HashingTF.transform() 将 words 列转换为特征向量,添加带有向量新列的 DataFrame
  3. 由于 LogisticRegressionEstimatorPipeline 首先调用 LogisticRegression.fit() 产生 LogisticRegressionModel,然后调用其 transform() 方法将 DataFrame 转换为另一个 DataFrame

Pipeline 本身是一个 Estimator。运行 Pipelinefit() 方法后,它会产生一个 PipelineModel,即一个 TransformerPipelineModel 用于测试阶段,原始中的所有 Estimator 已变为 Transformer

细节

  • DAG PipelinePipeline 被指定为一个有序数组。只要数据流图形成有向无环图(DAG),就可以创建非线性的 Pipeline。当前基于每个阶段的输入和输出列名称隐式指定该图。
  • 运行时检查PipelinePipelineModel 在实际运行之前会进行运行时检查,使用 DataFrame 模式完成类型检查。
  • 唯一的管道阶段Pipeline 的阶段应该是唯一的实例。例如,同一个 myHashingTF 实例不应插入两次,但不同的实例 myHashingTF1myHashingTF2(都属于 HashingTF 类型)可以放入同一个 Pipeline 中。

参数

MLlib EstimatorTransformer 使用统一的 API 来指定参数。

  • Param 是具有独立文件的命名参数
  • ParamMap 是一组(参数,值)对

将参数传递给算法的主要方式有两种:

  1. 设置实例的参数:例如 lr.setMaxIter(10),使 lr.fit() 最多使用 10 次迭代
  2. 传递 ParamMapfit()transform()ParamMap 中的任何参数将覆盖先前通过 setter 方法指定的参数

参数属于 EstimatorTransformer 的特定实例。例如,如果两个 LogisticRegression 实例 lr1lr2,则可以构建 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,包含 idtexts 列:

idtexts
0Array(“a”, “b”, “c”)
1Array(“a”, “b”, “b”, “c”, “a”)

CountVectorizer.fit() 产生一个 CountVectorizerModel,词汇表为 (a, b, c)。转换后的输出列 vector 包含:

idtextsvector
0Array(“a”, “b”, “c”)(3,[0,1,2],[1.0,1.0,1.0])
1Array(“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 训练后生成 CountVectorizerModelCountVectorizerModel 可以 loadsave。带有 Model 后缀的类也可以设置 setInputColsetOutputCol

FeatureHasher

FeatureHasher 将输入的多个特征通过 hash 函数进行向量化,返回结果为稀疏向量,稀疏向量中元素的值一般为整型数值,可以用于计算相似度。

假设有一个 DataFrame,包含 4 个输入列 realboolstringNumstring

realboolstringNumstring
2.2true1foo
3.3false2bar
4.4false3baz
5.5false4foo

FeatureHasher.transform 的输出:

realboolstringNumstringfeatures
2.2true1foo(262144,[51871, 63643,174475,253195],[1.0,1.0,2.2,1.0])
3.3false2bar(262144,[6031, 80619,140467,174475],[1.0,1.0,1.0,3.3])
4.4false3baz(262144,[24279,140467,174475,196810],[1.0,1.0,4.4,1.0])
5.5false4foo(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:

idraw
0[I, saw, the, red, baloon]
1[Mary, had, a, little, lamb]

应用 StopWordsRemover 后:

idrawfiltered
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 + y
  • setDegree(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:

idcategory
0a
1b
2c
3a
4a
5c

应用 StringIndexer 后:

idcategorycategoryIndex
0a0.0
1b2.0
2c1.0
3a0.0
4a0.0
5c1.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:

idcategoryIndex
00.0
12.0
21.0
30.0
40.0
51.0

应用 IndexToString 后:

idcategoryIndexoriginalCategory
00.0a
12.0b
21.0c
30.0a
40.0a
51.0c
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 中保存的元数据信息,用于恢复 categoryIndexoriginalCategory 的映射关系

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 编码。例如对于特征类别 123,对应 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

对多个向量取笛卡尔积,计算元素两两乘积。

假设输入列 id1vec1vec2Interaction 输出列 interactedCol

id1vec1vec2interactedCol
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):特征值将按照标准差进行缩放,确保方差等于 1
  • setWithMean(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:

idv1v2
01.03.0
22.05.0

使用 SQLTransformer 执行语句 "SELECT *, (v1 + v2) AS v3, (v1 * v2) AS v4 FROM __THIS__" 后:

idv1v2v3v4
01.03.04.03.0
22.05.07.010.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:

idhourmobileuserFeaturesclicked
0181.0[0.0, 10.0, 0.5]1.0

hourmobileuserFeatures 组合成 features 列后:

idhourmobileuserFeaturesclickedfeatures
0181.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,用户必须设置 inputColsize 参数。将此转换器应用于 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:

idhour
018.0
119.0
28.0
35.0
42.2

设置 numBuckets = 3 后:

idhourresult
018.02.0
119.02.0
28.01.0
35.01.0
42.20.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

用于填充数据中的缺失值,可以指定使用均值或中位数进行填充。

输入列应为 DoubleTypeFloatType。当前 Imputer 不支持分类特征,并且可能为包含分类特征的列创建不正确的值。Imputer 可以通过 .setMissingValue 填充除 NaN 以外的自定义值。例如 .setMissingValue(0) 将填充所有出现的 0

:输入列中的所有 null 值都被视为缺失,因此也进行了插补。

假设有以下 DataFrame:

ab
1.0Double.NaN
2.0Double.NaN
Double.NaN3.0
4.04.0
5.05.0

a 和列 b 的代理值分别为 3.04.0。转换后:

about_aout_b
1.0Double.NaN1.04.0
2.0Double.NaN2.04.0
Double.NaN3.03.03.0
4.04.04.04.0
5.05.05.05.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) 选择后两列:

userFeaturesfeatures
[0.0, 10.0, 0.5][10.0, 0.5]

如果输入属性为 ["f1", "f2", "f3"],使用 setNames("f2", "f3") 选择:

userFeaturesfeatures
[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 的元素
  • setIndicessetNames 同时使用时,取并集,保留两者的结果

RFormula

可以将 R 语言风格的公式应用于 DataFrame 数据。

基础操作

符号说明
~分离目标变量和项(terms)
+连接项,"+ 0" 表示移除截距
-移除一个项,"- 1" 表示移除截距
:交互项(数值乘法或二值化分类值的交互)
.除目标变量外的所有列

假设 ab 是 Double 列:

  • y ~ a + b → 模型 y ~ w0 + w1*a + w2*b,其中 w0 是截距,w1w2 是系数
  • y ~ a + b + a:b - 1 → 模型 y ~ w1*a + w2*b + w3*a*b,其中 w1w2w3 是系数

RFormula 产生一个特征向量列和一个 Double 或 String 类型的标签列。字符串列会先通过 StringIndexer 转换,然后进行独热编码。

stringOrderType 控制编码方式(假设字符串特征包含 {'b', 'a', 'b', 'a', 'c', 'b'}):

stringOrderTypeStringIndexer 映射到 0 的类别RFormula 丢弃的类别
frequencyDesc最频繁类别(b最不频繁类别(c
frequencyAsc最不频繁类别(c最频繁类别(b
alphabetDesc字母序最后的类别(c字母序最前的类别(a
alphabetAsc字母序最前的类别(a字母序最后的类别(c

假设有以下 DataFrame:

idcountryhourclicked
7”US”181.0
8”CA”120.0
9”NZ”150.0

使用公式 clicked ~ country + hour

idcountryhourclickedfeatureslabel
7”US”181.0[0.0, 0.0, 18.0]1.0
8”CA”120.0[0.0, 1.0, 12.0]0.0
9”NZ”150.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:

idfeaturesclicked
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

idfeaturesclickedselectedFeatures
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)设置损失函数类型:autobinomial(二分类)、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 可获取的信息

属性/方法说明
areaUnderROCROC 曲线下面积,0 到 1 之间,越接近 1 越好
rocROC 曲线,包含不同阈值下的真正例比率和假正例比率的 DataFrame
prPR(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)设置分割特征时的内存限制

setImpuritygini 不纯度值范围 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)设置特征选择策略:autoallsqrtonethird
setImpurity(value: String)设置不纯度度量:ginientropy
setSeed(value: Long)设置随机种子
setMaxBins(value: Int)设置最大分箱数量
setCacheNodeIds(value: Boolean)设置是否缓存节点 ID,加速特征重要性计算

setFeatureSubsetStrategyauto 根据问题类型自动选择;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)设置特征选择策略:allsqrtonethird
setImpurity(value: String)设置不纯度度量:ginientropy

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-bfgsgd

setLayers:以 [4, 5, 4, 3] 为例,第一个参数 4 为输入特征维度,最后一个参数 3 为输出类别数量,中间 54 为隐藏层大小。如果输入为 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)设置评估度量:accuracyf1weightedPrecisionweightedRecallweightedF1logloss
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 接口。

可用分布族

分布族返回类型支持的链接函数
GaussianContinuousIdentity*、Log、Inverse
BinomialBinaryLogit*、Probit、CLogLog
PoissonCountLog*、Identity、Sqrt
GammaContinuousInverse*、Identity、Log
TweedieZero-inflated continuousPower 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)设置链接函数:identityloginverselogitprobit
setFamily(value: String)设置分布族:gaussianbinomialpoissongammatweedie
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:调整方差的幂次数 pp=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 实现了带弹性网正则化的线性回归和逻辑回归。

决策树详解

决策树及其集成方法是分类和回归任务中最流行的方法之一。决策树被广泛使用,因为它们易于解释、能处理分类特征、扩展到多分类设置、不需要特征缩放、能够捕捉非线性和特征交互。

输入输出列

输入列

参数名类型默认值说明
labelColDouble"label"要预测的标签
featuresColVector"features"特征向量

输出列

参数名类型默认值说明备注
predictionColDouble"prediction"预测标签
rawPredictionColVector"rawPrediction"长度等于类别数的向量仅分类
probabilityColVector"probability"归一化为多项分布的 rawPrediction仅分类
varianceColDouble预测的有偏样本方差仅回归

与原始 MLlib 决策树 API 的主要区别:

  • 支持 ML Pipelines
  • 区分分类与回归的决策树
  • 使用 DataFrame 元数据区分连续和分类特征

树集成方法

DataFrame API 支持两种主要的树集成算法:随机森林(Random Forests)梯度提升树(GBTs),两者都使用决策树作为基础模型。

与原始 MLlib 集成 API 的主要区别:

  • 支持 DataFrames 和 ML Pipelines
  • 区分分类与回归
  • 使用 DataFrame 元数据区分连续和分类特征
  • 随机森林更多功能:特征重要性估计、分类中每个类别的预测概率

随机森林输入输出列

输入列

参数名类型默认值说明
labelColDouble"label"要预测的标签
featuresColVector"features"特征向量

输出列(预测)

参数名类型默认值说明备注
predictionColDouble"prediction"预测标签
rawPredictionColVector"rawPrediction"长度等于类别数的向量仅分类
probabilityColVector"probability"归一化为多项分布的 rawPrediction仅分类

GBT 输入输出列

注意GBTClassifier 当前仅支持二元标签。

参数名类型默认值说明备注
labelColDouble"label"要预测的标签
featuresColVector"features"特征向量
predictionColDouble"prediction"预测标签

未来 GBTClassifier 也将输出 rawPredictionprobability 列。

聚类

K-means

K-means 是最常用的聚类算法之一,将数据点聚类到预定义数量的簇中。MLlib 实现包括 k-means++ 方法的并行化变体(称为 kmeans||)。

输入输出列

参数名类型默认值说明
featuresColVector"features"特征向量
predictionColInt"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 作为同时支持 EMLDAOptimizerOnlineLDAOptimizerEstimator 实现,生成 LDAModelEMLDAOptimizer 生成的 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)算法来学习最大似然模型。

输入输出列

参数名类型默认值说明
featuresColVector"features"特征向量
predictionColInt"prediction"预测的簇中心
probabilityColVector"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 算法学习这些潜在因子。

参数说明

参数默认值说明
numBlocks10用户和物品分区并行计算的块数
rank10模型中潜在因子的数量
maxIter10最大迭代次数
regParam1.0ALS 中正则化参数
implicitPrefsfalse是否使用隐式反馈变体
alpha1.0隐式反馈变体的参数,控制偏好观察的基线置信度
nonnegativefalse是否对最小二乘使用非负约束

注意:基于 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(模型评估器)用于评价模型的类,包括 BinaryClassificationEvaluatorMulticlassClassificationEvaluatorRegressionEvaluatorClusteringEvaluator 等,具有 setMetricName 方法设置评估指标(如 AUC、RMSE 等)
Estimator(模型估计器)用于训练机器学习模型的抽象类,通常不具有 setMetricName 方法。例如 LogisticRegressionRandomForestClassifierKMeans
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设置要评估的机器学习模型(如 RandomForestLogisticRegressionPipeline 等)
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

CrossValidatorTrainValidationSplit
评估次数K 次(每个参数组合)1 次(每个参数组合)
可靠性依赖数据集大小
计算代价高(暴力穷举)较低
适用场景数据集不是特别大大数据集,快速迭代

Pipeline Stage 模型自定义

Transformer 自定义

Transformer 对数据做转换操作,不涉及模型训练,一般为中间结果,不需要保存到本地。效果为在原有 DataFrame 上增加一列。

创建步骤

  1. 继承 Transformer 抽象类
  2. [可选] 实现 setInputColsetOutputCol 方法
  3. 实现 transformSchema 方法
  4. 实现 transform 方法
  5. [可选] 实现可读写
  6. 通过 Pipeline 反射调用

步骤 1:[可选] 实现 setInputCol 和 setOutputCol 方法

初始化 inputColoutputCol 属性(类型为 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)。

创建步骤

  1. 构造通用参数 Trait
  2. 创建自定义 EstimatorModel 类(继承 Model、混入参数 Trait,重写 transformSchematransform
  3. 创建自定义 Estimator 类(继承 Estimator、混入参数 Trait,重写 transformSchemafit
  4. 实现可读写
  5. 通过 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 源码中 DefaultParamsWriterDefaultParamsReader 的方法是私有的,自定义类无法直接通过继承获得。需要从源码 spark.ml.util.ReadWrite.scala 中提取相关内容,实现 MyDefaultParamsWriterMyDefaultParamsReader 两个 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 + yaxpy(a: Double, x: Vector, y: Vector): Unit

Matrix(矩阵接口)

MatrixDenseMatrixSparseMatrix 的父类,用于访问和操作矩阵元素。

方法说明
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.lengthrowIndices 长度必须与 data 一致。

方法同 DenseMatrix,不同之处在于 colIter/rowIter 返回的是 SparseVector

Vector(向量接口)

VectorDenseVectorSparseVector 的父类。

方法说明
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())