Apache Spark MLlib 基础统计指南:RDD-based API 的六大统计能力详解
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
本文以 Apache Spark MLlib 的 RDD-based API 为对象,系统讲解六大基础统计能力:列汇总统计(colStats)、相关性分析(Pearson/Spearman)、分层抽样(sampleByKey/sampleByKeyExact)、假设检验(卡方检验与 Kolmogorov-Smirnov 检验)、随机数据生成与核密度估计。本文内容源自 mllib-statistics.md 官方文档,并结合仓库中的示例源码与核心实现深入展开,帮助读者在数据探索、特征筛选、A/B 测试与算法原型验证等场景中直接落地使用。
一、概览:RDD-based 统计 API 的定位
在 Spark MLlib 中,统计 API 主要分布在两个层级:基于 RDD 的org.apache.spark.mllib.stat.Statistics(Scala/Java)与pyspark.mllib.stat.Statistics(Python),以及基于 DataFrame 的新版 API。本文聚焦 RDD-based API,它直接面向RDD[Vector]、RDD[Double]、RDD[LabeledPoint]等原始分布式数据集,是早期 Spark 应用中最常用的一批统计工具。
RDD-based 统计 API 的入口是 Statistics.scala 中的object Statistics(@Since("1.1.0")),它聚合了colStats、corr、chiSqTest、kolmogorovSmirnovTest等核心方法;Python 侧对应 pyspark.mllib.stat.Statistics。下文按文档的六个板块逐一展开。
二、汇总统计(Summary Statistics)
2.1 API 与返回对象
对于RDD[Vector],通过Statistics.colStats计算逐列汇总统计,返回 MultivariateStatisticalSummary 实例,其中包含:
- count:样本总数(行数);
- max / min:每列的最大值 / 最小值;
- mean:每列均值;
- variance:每列方差;
- numNonzeros:每列非零元素个数。
从源码看,colStats(X: RDD[Vector])的实际实现是new RowMatrix(X).computeColumnSummaryStatistics()(见 Statistics.scala),即把 RDD 包装为分布式行矩阵后执行列统计计算,底层通过 treeAggregate 进行分区聚合,具备良好的分布式可扩展性。
2.2 Python 示例
完整示例见 summary_statistics_example.py:
import numpy as np from pyspark import SparkContext from pyspark.mllib.stat import Statistics sc = SparkContext(appName="SummaryStatisticsExample") # an RDD of Vectors(每行是一个样本,每列是一个维度) mat = sc.parallelize( [np.array([1.0, 10.0, 100.0]), np.array([2.0, 20.0, 200.0]), np.array([3.0, 30.0, 300.0])] ) # Compute column summary statistics. summary = Statistics.colStats(mat) print(summary.mean()) # a dense vector containing the mean value for each column print(summary.variance()) # column-wise variance print(summary.numNonzeros()) # number of nonzeros in each column sc.stop()2.3 Scala / Java 示例
Scala 完整示例见 SummaryStatisticsExample.scala:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.mllib.linalg.Vectors import org.apache.spark.mllib.stat.{MultivariateStatisticalSummary, Statistics} val conf = new SparkConf().setAppName("SummaryStatisticsExample") val sc = new SparkContext(conf) val observations = sc.parallelize( Seq( Vectors.dense(1.0, 10.0, 100.0), Vectors.dense(2.0, 20.0, 200.0), Vectors.dense(3.0, 30.0, 300.0) ) ) // Compute column summary statistics. val summary: MultivariateStatisticalSummary = Statistics.colStats(observations) println(summary.mean) // a dense vector containing the mean value for each column println(summary.variance) // column-wise variance println(summary.numNonzeros) // number of nonzeros in each column sc.stop()Java 对应示例为 JavaSummaryStatisticsExample.java,使用JavaRDD<Vector>调用Statistics.colStats。注意:MultivariateStatisticalSummary返回的mean、variance等都是稠密向量(DenseVector),逐列下标与输入矩阵的列一一对应,便于在特征标准化、数据质量探查等场景中直接使用。
三、相关性分析(Correlations)
计算两组数据序列之间的相关性是统计学中最常见的操作之一。spark.mllib目前支持Pearson 相关系数与Spearman 秩相关系数两种方法,并且可以计算多序列之间的两两相关矩阵。
3.1 输入与输出规则
根据输入类型的不同,Statistics.corr的返回不同:
| 输入类型 | 输出 |
|---|---|
两个RDD[Double](Python)/JavaDoubleRDD(Java) | 相关系数Double |
一个RDD[Vector](Python)/JavaRDD<Vector>(Java) | 相关性矩阵Matrix |
方法参数method支持"pearson"(默认)与"spearman",不传时默认使用 Pearson 方法。
3.2 Python 示例
完整示例见 correlations_example.py:
import numpy as np from pyspark import SparkContext from pyspark.mllib.stat import Statistics sc = SparkContext(appName="CorrelationsExample") seriesX = sc.parallelize([1.0, 2.0, 3.0, 3.0, 5.0]) # a series # seriesY must have the same number of partitions and cardinality as seriesX seriesY = sc.parallelize([11.0, 22.0, 33.0, 33.0, 555.0]) # Compute the correlation using Pearson's method. Enter "spearman" for Spearman's method. # If a method is not specified, Pearson's method will be used by default. print("Correlation is: " + str(Statistics.corr(seriesX, seriesY, method="pearson"))) data = sc.parallelize( [np.array([1.0, 10.0, 100.0]), np.array([2.0, 20.0, 200.0]), np.array([5.0, 33.0, 366.0])] ) # an RDD of Vectors # calculate the correlation matrix using Pearson's method. Use "spearman" for Spearman's method. # If a method is not specified, Pearson's method will be used by default. print(Statistics.corr(data, method="pearson")) sc.stop()3.3 Scala 示例与实现细节
Scala 完整示例见 CorrelationsExample.scala:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.mllib.linalg._ import org.apache.spark.mllib.stat.Statistics import org.apache.spark.rdd.RDD val conf = new SparkConf().setAppName("CorrelationsExample") val sc = new SparkContext(conf) // 两个同长度、同分区数的 Double 序列 val seriesX: RDD[Double] = sc.parallelize(Array(1.0, 2.0, 3.0, 3.0, 5.0)) val seriesY: RDD[Double] = sc.parallelize(Array(11.0, 22.0, 33.0, 33.0, 555.0)) val correlation: Double = Statistics.corr(seriesX, seriesY, "pearson") println(s"Correlation is: $correlation") // 注意:每个 Vector 是一行(一个样本),而不是一列 val data: RDD[Vector] = sc.parallelize( Seq( Vectors.dense(1.0, 10.0, 100.0), Vectors.dense(2.0, 20.0, 200.0), Vectors.dense(5.0, 33.0, 366.0)) ) // 计算相关矩阵 val correlMatrix: Matrix = Statistics.corr(data, "pearson") println(correlMatrix.toString) sc.stop()结合 Statistics.scala 的源码注释,有几点实现细节值得注意:
- Spearman 的代价更高:Spearman 是秩相关,需要对每一列构造
RDD[Double]并排序以获取秩,再将各列 join 回RDD[Vector],过程相当昂贵。官方建议:调用corr且method = "spearman"之前先对输入 RDD 做 cache,避免公共血缘被重复计算。 - 零方差/零协方差边界:两序列相关时,若任一向量方差为 0,返回
NaN;相关矩阵中协方差为 0 的列对会产生NaN条目(见 Statistics.scala 与corr(x, y)的注释)。 - 分区一致性约束:两个
RDD[Double]必须具有相同的分区数与每个分区内相同的元素数,相关计算依赖逐分区配对。 - 列方向说明:
RDD[Vector]求相关矩阵时,每个Vector是一行(样本),相关矩阵的行列对应的是特征列之间的相关性,矩阵为对称方阵。
四、分层抽样(Stratified Sampling)
4.1 概念与两种方法对比
与其他位于spark.mllib包内的统计函数不同,分层抽样方法sampleByKey与sampleByKeyExact是定义在key-value 对 RDD上的算子(Scala 侧位于RDD[(K, V)]的隐式类 PairRDDFunctions,Java 侧位于 JavaPairRDD)。分层抽样中,key 可视为"标签",value 为具体属性:例如 key 是男/女或文档 ID,value 是人群年龄列表或文档中的单词列表。
两种方法的取舍:
sampleByKey:对每个观测以概率决定是否抽取("掷硬币"方式),只需一次数据遍历,提供的是期望样本量;sampleByKeyExact:需要比sampleByKey的逐层简单随机抽样显著更多的资源,但能以99.99% 的置信度给出精确的抽样数量:无放回抽样需要额外一次数据遍历,有放回抽样需要额外两次遍历;- 注意:
sampleByKeyExact当前不支持 Python(pyspark.RDD.sampleByKey可用,但无 exact 版本)。
4.2 抽样公式
两种方法的目标都是对每个 keyk精确抽取约 $\lceil f_k \cdot n_k \rceil$ 个样本(∀ k ∈ K),其中:
- $f_k$:为 key
k指定的期望抽样比例; - $n_k$:key
k对应的 key-value 对总数; - $K$:所有 key 的集合。
4.3 Python 示例(仅近似抽样)
完整示例见 stratified_sampling_example.py:
from pyspark import SparkContext sc = SparkContext(appName="StratifiedSamplingExample") # an RDD of any key value pairs data = sc.parallelize([(1, 'a'), (1, 'b'), (2, 'c'), (2, 'd'), (2, 'e'), (3, 'f')]) # specify the exact fraction desired from each key as a dictionary fractions = {1: 0.1, 2: 0.6, 3: 0.3} approxSample = data.sampleByKey(False, fractions) # $example off$ for each in approxSample.collect(): print(each) sc.stop()4.4 Scala 示例(近似 + 精确)
完整示例见 StratifiedSamplingExample.scala:
import org.apache.spark.{SparkConf, SparkContext} val conf = new SparkConf().setAppName("StratifiedSamplingExample") val sc = new SparkContext(conf) // an RDD[(K, V)] of any key value pairs val data = sc.parallelize( Seq((1, 'a'), (1, 'b'), (2, 'c'), (2, 'd'), (2, 'e'), (3, 'f'))) // specify the exact fraction desired from each key val fractions = Map(1 -> 0.1, 2 -> 0.6, 3 -> 0.3) // Get an approximate sample from each stratum val approxSample = data.sampleByKey(withReplacement = false, fractions = fractions) // Get an exact sample from each stratum val exactSample = data.sampleByKeyExact(withReplacement = false, fractions = fractions) println(s"approxSample size is ${approxSample.collect().size}") approxSample.collect().foreach(println) println(s"exactSample its size is ${exactSample.collect().size}") exactSample.collect().foreach(println) sc.stop()4.5 实现要点(从源码结构推断)
sampleByKeyExact的实现核心是先采样估算每个 key 的分布,再计算达到目标样本量所需的精确采样比例,最后再执行一次(无放回)或两次(有放回)补充遍历,因此能以 99.99% 置信度保证每个 key 的样本量精确等于 $\lceil f_k \cdot n_k \rceil$;- 精确抽样在 key 数量庞大或数据倾斜严重时,额外遍历与分布估算的开销会明显高于
sampleByKey,因此在可以接受近似样本量的场景下应优先选择sampleByKey。
五、假设检验(Hypothesis Testing)
假设检验用于判断某个结果是否具有统计显著性——即结果到底是真实效应还是偶然发生。spark.mllib当前支持Pearson 卡方(χ²)检验与单样本双边 Kolmogorov-Smirnov(KS)检验。
5.1 卡方检验:由输入类型决定检验方式
卡方检验的输入数据类型决定了执行哪种检验:
| 输入类型 | 执行的检验 |
|---|---|
Vector | 拟合优度检验(goodness of fit):观测频率是否与期望分布一致 |
Matrix(列联表) | 独立性检验(independence):两个分类变量是否相互独立 |
RDD[LabeledPoint] | 特征选择:对每个特征分别与标签做独立性检验,返回ChiSqTestResult数组 |
关于Vector输入的实现细节,源码注释说明(见 Statistics.scala):若未提供期望分布向量,检验默认对均匀分布进行;若提供了expected向量,当expected总和与observed总和不一致时,expected会被重新缩放。
5.2 Python 示例
完整示例见 hypothesis_testing_example.py:
from pyspark import SparkContext from pyspark.mllib.linalg import Matrices, Vectors from pyspark.mllib.regression import LabeledPoint from pyspark.mllib.stat import Statistics sc = SparkContext(appName="HypothesisTestingExample") # a vector composed of the frequencies of events vec = Vectors.dense(0.1, 0.15, 0.2, 0.3, 0.25) # compute the goodness of fit. If a second vector to test against # is not supplied as a parameter, the test runs against a uniform distribution. goodnessOfFitTestResult = Statistics.chiSqTest(vec) print("%s\n" % goodnessOfFitTestResult) # a contingency matrix(3 行 2 列的列联表) mat = Matrices.dense(3, 2, [1.0, 3.0, 5.0, 2.0, 4.0, 6.0]) # conduct Pearson's independence test on the input contingency matrix independenceTestResult = Statistics.chiSqTest(mat) print("%s\n" % independenceTestResult) # LabeledPoint(label, feature),用于特征选择 obs = sc.parallelize( [LabeledPoint(1.0, [1.0, 0.0, 3.0]), LabeledPoint(1.0, [1.0, 2.0, 0.0]), LabeledPoint(1.0, [-1.0, 0.0, -0.5])] ) # 由 RDD[LabeledPoint] 构造列联表,对每个特征与标签做独立性检验, # 返回一个数组,元素是每个特征的 ChiSquaredTestResult featureTestResults = Statistics.chiSqTest(obs) for i, result in enumerate(featureTestResults): print("Column %d:\n%s" % (i + 1, result)) sc.stop()5.3 Scala 示例
完整示例见 HypothesisTestingExample.scala:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.mllib.linalg._ import org.apache.spark.mllib.regression.LabeledPoint import org.apache.spark.mllib.stat.Statistics import org.apache.spark.mllib.stat.test.ChiSqTestResult import org.apache.spark.rdd.RDD val conf = new SparkConf().setAppName("HypothesisTestingExample") val sc = new SparkContext(conf) // 事件频率向量:拟合优度检验(不提供期望分布则默认对均匀分布检验) val vec: Vector = Vectors.dense(0.1, 0.15, 0.2, 0.3, 0.25) val goodnessOfFitTestResult = Statistics.chiSqTest(vec) println(s"$goodnessOfFitTestResult\n") // 3x2 列联表:独立性检验 val mat: Matrix = Matrices.dense(3, 2, Array(1.0, 3.0, 5.0, 2.0, 4.0, 6.0)) val independenceTestResult = Statistics.chiSqTest(mat) println(s"$independenceTestResult\n") // (label, feature) 对:逐特征独立性检验,用于特征选择 val obs: RDD[LabeledPoint] = sc.parallelize( Seq( LabeledPoint(1.0, Vectors.dense(1.0, 0.0, 3.0)), LabeledPoint(1.0, Vectors.dense(1.0, 2.0, 0.0)), LabeledPoint(-1.0, Vectors.dense(-1.0, 0.0, -0.5)) ) ) val featureTestResults: Array[ChiSqTestResult] = Statistics.chiSqTest(obs) featureTestResults.zipWithIndex.foreach { case (k, v) => println(s"Column ${(v + 1)} :") println(k) } sc.stop()Java 侧完整示例为 JavaHypothesisTestingExample.java,检验结果类型可参考 ChiSqTestResult。ChiSqTestResult的字符串输出包含:p 值(p-value)、自由度(degrees of freedom)、检验统计量(test statistic)、使用的检验方法以及原假设(null hypothesis)描述。判读规则:当 p 值小于显著性水平(如 0.05)时,拒绝原假设,认为差异具有统计显著性。
5.4 Kolmogorov-Smirnov(KS)检验
spark.mllib提供1 样本、双边的 KS 检验,用于检验样本是否来自某个理论分布(分布相等性检验)。用户有两种指定理论分布的方式:
- 提供理论分布的名称及其参数——当前仅支持正态分布(
distName = "norm")。若测试正态分布但未提供分布参数,检验会初始化为标准正态分布N(0, 1)并打印相应日志; - 提供一个计算累积分布函数(CDF)的函数(Scala 支持传入 lambda,Python API 不提供该能力)。
Python 完整示例见 hypothesis_testing_kolmogorov_smirnov_test_example.py:
from pyspark import SparkContext from pyspark.mllib.stat import Statistics sc = SparkContext(appName="HypothesisTestingKolmogorovSmirnovTestExample") parallelData = sc.parallelize([0.1, 0.15, 0.2, 0.3, 0.25]) # run a KS test for the sample versus a standard normal distribution testResult = Statistics.kolmogorovSmirnovTest(parallelData, "norm", 0, 1) # 结果包括 p 值、检验统计量与原假设; # 若 p 值表明显著,则可拒绝原假设。 print(testResult) sc.stop()Scala 侧对应示例为 HypothesisTestingKolmogorovSmirnovTestExample.scala,除指定分布名与参数外,还可传入自定义 CDF lambda:
// 示例:测试样本是否来自正态分布 N(0, 1) val testResult = Statistics.kolmogorovSmirnovTest(parallelData, "norm", 0, 1) // 或传入自定义 CDF 函数 val myCDF = (x: Double) => 1.0 - math.exp(-x) // 以指数分布 CDF 为例六、流式显著性检验(Streaming Significance Testing)
为支持A/B 测试等实时场景,spark.mllib提供了假设检验的在线(流式)实现。检验作用在 Spark Streaming 的DStream[(Boolean, Double)]上:每个元组第一个元素表示对照组(false)或实验组(true),第二个元素是观测值。
6.1 参数说明
流式显著性检验支持以下参数:
peacePeriod:流开始时要忽略的初始数据点数量,用于缓解新奇效应(novelty effects);windowSize:进行假设检验所覆盖的过去批次数量。设为0表示累积处理——使用全部历史批次。
6.2 使用方式
Scala 侧通过 StreamingTest 提供流式假设检验,完整示例见 StreamingTestExample.scala:
import org.apache.spark.mllib.stat.test.StreamingTest // 创建 StreamingTest,设置静默期与窗口大小 val test = new StreamingTest() .setPeacePeriod(0) .setWindowSize(0) .setTestMethod("welch") // 支持的方法见下 // 输入 DStream[(Boolean, Double)],对每个批次输出检验结果 val dstream: DStream[(Boolean, Double)] = ... dstream.foreachRDD { rdd => val result = rdd.map(testResult => testResult.pValue) // 处理每批次的 p 值结果 }Java 侧对应示例为 JavaStreamingTestExample.java,使用JavaDStream<Tuple2<Boolean, Double>>作为输入。StreamingTest的检验方法通过setTestMethod配置(可参考 StreamingTest.scala 中的StreamingTestMethod支持集合,如 Welch 检验与 Student 检验),输出为带 p 值的检验结果流,便于实时监控对照组与实验组的差异显著性。
七、随机数据生成(Random Data Generation)
随机数据生成对于随机化算法、原型验证和性能测试非常有用。spark.mllib支持生成服从指定分布的 i.i.d.(独立同分布)随机 RDD,支持均匀分布(uniform)、标准正态分布(standard normal)与泊松分布(Poisson)。
7.1 Python 示例
Python 侧工厂方法位于 RandomRDDs,以下示例生成 100 万个服从N(0, 1)标准正态分布的值(均匀分布在 10 个分区),再映射为N(1, 4):
from pyspark.mllib.random import RandomRDDs sc = ... # SparkContext # Generate a random double RDD that contains 1 million i.i.d. values drawn from the # standard normal distribution `N(0, 1)`, evenly distributed in 10 partitions. u = RandomRDDs.normalRDD(sc, 1000000, 10) # Apply a transform to get a random double RDD following `N(1, 4)`. v = u.map(lambda x: 1.0 + 2.0 * x)7.2 Scala 示例
Scala 侧工厂方法位于 RandomRDDs,完整示例见 RandomRDDGeneration.scala:
import org.apache.spark.SparkContext import org.apache.spark.mllib.random.RandomRDDs._ val sc: SparkContext = ... // Generate a random double RDD that contains 1 million i.i.d. values drawn from the // standard normal distribution `N(0, 1)`, evenly distributed in 10 partitions. val u = normalRDD(sc, 1000000L, 10) // Apply a transform to get a random double RDD following `N(1, 4)`. val v = u.map(x => 1.0 + 2.0 * x)7.3 Java 示例
Java 侧使用normalJavaRDD:
import org.apache.spark.SparkContext; import org.apache.spark.api.JavaDoubleRDD; import static org.apache.spark.mllib.random.RandomRDDs.*; JavaSparkContext jsc = ... // Generate a random double RDD that contains 1 million i.i.d. values drawn from the // standard normal distribution `N(0, 1)`, evenly distributed in 10 partitions. JavaDoubleRDD u = normalJavaRDD(jsc, 1000000L, 10); // Apply a transform to get a random double RDD following `N(1, 4)`. JavaDoubleRDD v = u.mapToDouble(x -> 1.0 + 2.0 * x);7.4 工厂方法速查(从源码结构推断)
RandomRDDs 为每种分布提供 double 型与 vector 型两个变体,方法命名规律为<dist>RDD与<dist>VectorRDD,如:
uniformRDD/uniformVectorRDD:均匀分布;normalRDD/normalVectorRDD:标准正态分布;poissonRDD/poissonVectorRDD:泊松分布(需传入mean参数)。
每个工厂方法的典型签名为(sc, size, numPartitions, seed),seed不传时使用随机种子,传入固定种子可保证实验结果可复现。该类随机 RDD 常用于性能基准测试的数据生成(本仓库mllib/benchmarks与core/benchmarks下的基准测试数据即由此类方法生成)与无真实数据集时的算法原型验证。
八、核密度估计(Kernel Density Estimation)
核密度估计(KDE)是一种无需对观测样本的潜在分布做任何假设即可可视化经验概率分布的技术。它计算随机变量概率密度函数(PDF)在给定评估点集合上的估计值——将经验分布在某一点的 PDF 表示为以每个样本为中心的正态分布 PDF 的均值(即高斯核平滑)。
8.1 Python 示例
Python 侧由 KernelDensity 提供,完整示例见 kernel_density_estimation_example.py:
from pyspark import SparkContext from pyspark.mllib.stat import KernelDensity sc = SparkContext(appName="KernelDensityEstimationExample") # an RDD of sample data data = sc.parallelize([1.0, 1.0, 1.0, 2.0, 3.0, 4.0, 5.0, 5.0, 6.0, 7.0, 8.0, 9.0, 9.0]) # Construct the density estimator with the sample data and a standard deviation for the Gaussian # kernels kd = KernelDensity() kd.setSample(data) kd.setBandwidth(3.0) # Find density estimates for the given values densities = kd.estimate([-1.0, 2.0, 5.0]) print(densities) sc.stop()8.2 Scala 示例与实现要点
Scala 完整示例见 KernelDensityEstimationExample.scala:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.mllib.stat.KernelDensity import org.apache.spark.rdd.RDD val conf = new SparkConf().setAppName("KernelDensityEstimationExample") val sc = new SparkContext(conf) // an RDD of sample data val data: RDD[Double] = sc.parallelize(Seq(1, 1, 1, 2, 3, 4, 5, 5, 6, 7, 8, 9, 9)) // Construct the density estimator with the sample data and a standard deviation // for the Gaussian kernels val kd = new KernelDensity() .setSample(data) .setBandwidth(3.0) // Find density estimates for the given values val densities = kd.estimate(Array(-1.0, 2.0, 5.0)) densities.foreach(println) sc.stop()从源码结构看,核心实现位于 KernelDensity.scala:setBandwidth设置高斯核的标准差(带宽),estimate在给定评估点上,把每个样本视为一个高斯核的均值点,将所有核在该点的密度贡献取平均得到估计值。带宽越小,估计越"尖锐"(更贴合样本);带宽越大,曲线越平滑。KDE 非常适合分布形态探查:在直方图对分布形状过于敏感时,KDE 输出的是平滑连续的概率密度曲线,可直接用于绘图与异常检测的基线建模。
九、总结与选择建议
| 功能 | 核心 API | 适用场景 |
|---|---|---|
| 列汇总统计 | Statistics.colStats→MultivariateStatisticalSummary | 数据质量探查、特征标准化前的均值/方差计算 |
| 相关性分析 | Statistics.corr(rddX, rddY, method)/Statistics.corr(rddVector, method) | 特征共线性诊断、变量关联度探索 |
| 分层抽样 | RDD.sampleByKey/RDD.sampleByKeyExact | 按标签分层采样、类别不平衡数据重采样(精确样本量选 exact 版本) |
| 假设检验 | Statistics.chiSqTest/Statistics.kolmogorovSmirnovTest | 分布拟合、独立性检验、特征选择 |
| 流式显著性检验 | StreamingTest(配合DStream[(Boolean, Double)]) | A/B 测试在线监控 |
| 随机数据生成 | RandomRDDs(uniform/normal/poisson) | 算法原型、性能基准测试 |
| 核密度估计 | KernelDensity(setSample+setBandwidth+estimate) | 经验分布可视化、无假设密度估计 |
使用建议(综合源码与文档给出的注意点):
- 计算Spearman 相关前先对输入 RDD 执行
cache(),避免血缘重算; - 两个
RDD[Double]序列求相关时,确保分区数与每分区元素数一致; - 可接受近似样本量时优先用
sampleByKey,需要精确样本量(如要求 $\lceil f_k \cdot n_k \rceil$)时用sampleByKeyExact(Python 不可用); - KS 检验目前只内建支持正态分布,其他理论分布需自行提供 CDF 函数(Python API 不支持 lambda 形式);
- 随机数据生成建议传入固定
seed,保证实验可复现。
读者可进一步参考仓库内的完整示例集合:examples/src/main/scala/org/apache/spark/examples/mllib、examples/src/main/python/mllib 与 examples/src/main/java/org/apache/spark/examples/mllib,以及 API 文档 mllib-guide.md,结合本文内容即可快速在 Spark 集群上开展统计分析与实验。
【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考