
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.StatisticsScala/Java与pyspark.mllib.stat.StatisticsPython以及基于 DataFrame 的新版 API。本文聚焦 RDD-based API它直接面向RDD[Vector]、RDD[Double]、RDD[LabeledPoint]等原始分布式数据集是早期 Spark 应用中最常用的一批统计工具。RDD-based 统计 API 的入口是 Statistics.scala 中的object StatisticsSince(1.1.0)它聚合了colStats、corr、chiSqTest、kolmogorovSmirnovTest等核心方法Python 侧对应 pyspark.mllib.stat.Statistics。下文按文档的六个板块逐一展开。二、汇总统计Summary Statistics2.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.pyimport numpy as np from pyspark import SparkContext from pyspark.mllib.stat import Statistics sc SparkContext(appNameSummaryStatisticsExample) # 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.scalaimport 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使用JavaRDDVector调用Statistics.colStats。注意MultivariateStatisticalSummary返回的mean、variance等都是稠密向量DenseVector逐列下标与输入矩阵的列一一对应便于在特征标准化、数据质量探查等场景中直接使用。三、相关性分析Correlations计算两组数据序列之间的相关性是统计学中最常见的操作之一。spark.mllib目前支持Pearson 相关系数与Spearman 秩相关系数两种方法并且可以计算多序列之间的两两相关矩阵。3.1 输入与输出规则根据输入类型的不同Statistics.corr的返回不同输入类型输出两个RDD[Double]Python/JavaDoubleRDDJava相关系数Double一个RDD[Vector]Python/JavaRDDVectorJava相关性矩阵Matrix方法参数method支持pearson默认与spearman不传时默认使用 Pearson 方法。3.2 Python 示例完整示例见 correlations_example.pyimport numpy as np from pyspark import SparkContext from pyspark.mllib.stat import Statistics sc SparkContext(appNameCorrelationsExample) 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 Pearsons method. Enter spearman for Spearmans method. # If a method is not specified, Pearsons method will be used by default. print(Correlation is: str(Statistics.corr(seriesX, seriesY, methodpearson))) 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 Pearsons method. Use spearman for Spearmans method. # If a method is not specified, Pearsons method will be used by default. print(Statistics.corr(data, methodpearson)) sc.stop()3.3 Scala 示例与实现细节Scala 完整示例见 CorrelationsExample.scalaimport 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(sCorrelation 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 Sampling4.1 概念与两种方法对比与其他位于spark.mllib包内的统计函数不同分层抽样方法sampleByKey与sampleByKeyExact是定义在key-value 对 RDD上的算子Scala 侧位于RDD[(K, V)]的隐式类 PairRDDFunctionsJava 侧位于 JavaPairRDD。分层抽样中key 可视为标签value 为具体属性例如 key 是男/女或文档 IDvalue 是人群年龄列表或文档中的单词列表。两种方法的取舍sampleByKey对每个观测以概率决定是否抽取掷硬币方式只需一次数据遍历提供的是期望样本量sampleByKeyExact需要比sampleByKey的逐层简单随机抽样显著更多的资源但能以99.99% 的置信度给出精确的抽样数量无放回抽样需要额外一次数据遍历有放回抽样需要额外两次遍历注意sampleByKeyExact当前不支持 Pythonpyspark.RDD.sampleByKey可用但无 exact 版本。4.2 抽样公式两种方法的目标都是对每个 keyk精确抽取约 $\lceil f_k \cdot n_k \rceil$ 个样本∀ k ∈ K其中$f_k$为 keyk指定的期望抽样比例$n_k$keyk对应的 key-value 对总数$K$所有 key 的集合。4.3 Python 示例仅近似抽样完整示例见 stratified_sampling_example.pyfrom pyspark import SparkContext sc SparkContext(appNameStratifiedSamplingExample) # 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.scalaimport 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(sapproxSample size is ${approxSample.collect().size}) approxSample.collect().foreach(println) println(sexactSample 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-SmirnovKS检验。5.1 卡方检验由输入类型决定检验方式卡方检验的输入数据类型决定了执行哪种检验输入类型执行的检验Vector拟合优度检验goodness of fit观测频率是否与期望分布一致Matrix列联表独立性检验independence两个分类变量是否相互独立RDD[LabeledPoint]特征选择对每个特征分别与标签做独立性检验返回ChiSqTestResult数组关于Vector输入的实现细节源码注释说明见 Statistics.scala若未提供期望分布向量检验默认对均匀分布进行若提供了expected向量当expected总和与observed总和不一致时expected会被重新缩放。5.2 Python 示例完整示例见 hypothesis_testing_example.pyfrom pyspark import SparkContext from pyspark.mllib.linalg import Matrices, Vectors from pyspark.mllib.regression import LabeledPoint from pyspark.mllib.stat import Statistics sc SparkContext(appNameHypothesisTestingExample) # 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 matrix3 行 2 列的列联表 mat Matrices.dense(3, 2, [1.0, 3.0, 5.0, 2.0, 4.0, 6.0]) # conduct Pearsons 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.scalaimport 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(sColumn ${(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-SmirnovKS检验spark.mllib提供1 样本、双边的 KS 检验用于检验样本是否来自某个理论分布分布相等性检验。用户有两种指定理论分布的方式提供理论分布的名称及其参数——当前仅支持正态分布distName norm。若测试正态分布但未提供分布参数检验会初始化为标准正态分布N(0, 1)并打印相应日志提供一个计算累积分布函数CDF的函数Scala 支持传入 lambdaPython API 不提供该能力。Python 完整示例见 hypothesis_testing_kolmogorov_smirnov_test_example.pyfrom pyspark import SparkContext from pyspark.mllib.stat import Statistics sc SparkContext(appNameHypothesisTestingKolmogorovSmirnovTestExample) 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 effectswindowSize进行假设检验所覆盖的过去批次数量。设为0表示累积处理——使用全部历史批次。6.2 使用方式Scala 侧通过 StreamingTest 提供流式假设检验完整示例见 StreamingTestExample.scalaimport 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使用JavaDStreamTuple2Boolean, 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.scalaimport 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 侧使用normalJavaRDDimport 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 型两个变体方法命名规律为distRDD与distVectorRDD如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.pyfrom pyspark import SparkContext from pyspark.mllib.stat import KernelDensity sc SparkContext(appNameKernelDensityEstimationExample) # 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.scalaimport 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.scalasetBandwidth设置高斯核的标准差带宽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 测试在线监控随机数据生成RandomRDDsuniform/normal/poisson算法原型、性能基准测试核密度估计KernelDensitysetSamplesetBandwidthestimate经验分布可视化、无假设密度估计使用建议综合源码与文档给出的注意点计算Spearman 相关前先对输入 RDD 执行cache()避免血缘重算两个RDD[Double]序列求相关时确保分区数与每分区元素数一致可接受近似样本量时优先用sampleByKey需要精确样本量如要求 $\lceil f_k \cdot n_k \rceil$时用sampleByKeyExactPython 不可用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),仅供参考