You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在分组PySpark DataFrame中推导皮尔逊相关系数的实现方案

解决方案

要实现按customer_name分组计算price与units的皮尔逊相关系数,可通过**分组后对每个子数据集使用pyspark.ml.stat.Correlation**完成,具体步骤如下:

步骤说明

  1. 导入依赖库:用VectorAssembler将price和units合并为特征向量(Correlation要求输入为向量列),用Correlation计算相关系数。
  2. 定义分组处理函数:对每个分组的子数据集,先转换为向量列,再计算皮尔逊相关系数矩阵,提取对应位置的系数值。
  3. 分组并应用处理函数:使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 19:52:43