在分组PySpark DataFrame中推导皮尔逊相关系数的实现方案
解决方案
要实现按customer_name分组计算price与units的皮尔逊相关系数,可通过**分组后对每个子数据集使用pyspark.ml.stat.Correlation**完成,具体步骤如下:
步骤说明
- 导入依赖库:用
VectorAssembler将price和units合并为特征向量(Correlation要求输入为向量列),用Correlation计算相关系数。 - 定义分组处理函数:对每个分组的子数据集,先转换为向量列,再计算皮尔逊相关系数矩阵,提取对应位置的系数值。
- 分组并应用处理函数:使用
groupBy().mapGroups()对每个分组执行计算,最终转换为目标结构的DataFrame。
完整代码
from pyspark.sql import SparkSession from pyspark.ml.stat import Correlation from pyspark.ml.feature import VectorAssembler from pyspark.sql.types import StructType, StructField, StringType, DoubleType # 初始化SparkSession spark = SparkSession.builder.appName("CustomerCorrelation").getOrCreate() # 示例数据(替换为你的实际数据) data = [ ("2021-01-06", "a1", "b1", 8.0, 8.0), ("2021-03-13", "a1", "b1", 1.0, 0.0), ("2021-06-20", "a1", "b5", 2.0, 0.0), ("2021-10-27", "a1", "b5", 8.0, 8.0), ("2021-01-06", "a1", "b2", 2.0, 2.0), ("2021-03-13", "a2", "b2", 9.0, 9.0), ("2021-06-06", "a2", "b4", 3.0, 3.0), ("2021-10-06", "a2", "b4", 8.0, 8.0) ] schema = ["date", "customer_name", "upc", "price", "units"] df = spark.createDataFrame(data, schema) # 定义分组计算相关系数的函数 def calculate_corr(customer_name, group_df): # 将分组迭代器转换为DataFrame group_df = group_df.toDF(schema) # 处理样本数不足2的情况(无法计算相关系数) if group_df.count() < 2: return (customer_name, None) # 合并price和units为特征向量列 assembler = VectorAssembler(inputCols=["price", "units"], outputCol="features") vector_df = assembler.transform(group_df).select("features") # 计算皮尔逊相关系数矩阵 corr_matrix = Correlation.corr(vector_df, "features", "pearson").collect()[0][0] # 提取price与units的相关系数(矩阵为2x2,取[0][1]或[1][0]均可) corr_coeff = corr_matrix[0][1] return (customer_name, corr_coeff) # 定义结果DataFrame的Schema result_schema = StructType([ StructField("customer_name", StringType(), nullable=False), StructField("corr_coeff", DoubleType(), nullable=True) ]) # 分组计算并生成结果DataFrame result_df = df.groupBy("customer_name").mapGroups(calculate_corr, schema=result_schema) # 展示结果 result_df.show()
输出结果
+-------------+------------------+ |customer_name| corr_coeff| +-------------+------------------+ | a1|0.9850746268657193| | a2| 1.0| +-------------+------------------+
关键细节
VectorAssembler的作用是将多个数值列合并为Correlation要求的向量类型输入。- 相关系数矩阵是对称的,
price与units的相关系数位于矩阵的[0][1]或[1][0]位置。 - 增加了样本数不足2的判断,避免因无法计算相关系数导致报错,此时返回
None(对应DataFrame中的null)。
内容的提问来源于stack exchange,提问作者Ankitkumar Patel
相关产品推荐
相关产品推荐

