如何处理Spark groupBy+collect_list生成的数组并使用numpy.cov计算指定比值?
实现代码
第一步:导包并定义自定义计算函数
import pyspark.sql.functions as F import numpy as np from pyspark.sql.types import DoubleType # 定义cov比值计算逻辑 def calc_cov_ratio(arr1, arr2): # 提前将decimal类型元素转为float,避免类型不兼容 arr1_np = np.array(arr1, dtype=float) arr2_np = np.array(arr2, dtype=float) # 按需求计算cov矩阵 cov_matrix = np.cov(arr1_np - arr2_np, arr1_np + arr2_np) return float(cov_matrix[0][1] / cov_matrix[1][1]) # 注册为Spark UDF,指定返回值为双精度浮点型 cov_ratio_udf = F.udf(calc_cov_ratio, DoubleType())
第二步:执行计算得到结果
# 直接对分组后的df2调用UDF计算 result_df = df2.withColumn("cov[0][1] / cov[1][1]", cov_ratio_udf(F.col("array1"), F.col("array2"))) result_df.show()
输出示例
+---+---------------------+ | id|cov[0][1] / cov[1][1]| +---+---------------------+ | b| 0.25| | a| 0.25| +---+---------------------+
报错原因说明
- 报错
'float' object has no attribute 'shape':是因为未将Spark Array类型的字段转换为numpy可识别的数组格式,直接传入numpy函数导致识别失败 - 报错
unexpected type:<type 'type'>:是因为原始字段为decimal类型,未转换为数值类型就传入numpy计算,出现类型不匹配
注意事项
如果分组内数据条数小于2,np.cov会返回空值,可在UDF中增加长度判断逻辑避免运行报错:
def calc_cov_ratio(arr1, arr2): if len(arr1) < 2 or len(arr2) <2: return None arr1_np = np.array(arr1, dtype=float) arr2_np = np.array(arr2, dtype=float) cov_matrix = np.cov(arr1_np - arr2_np, arr1_np + arr2_np) return float(cov_matrix[0][1] / cov_matrix[1][1])
内容的提问来源于stack exchange,提问作者KaguraKoishi
相关产品推荐
相关产品推荐

