基于PySpark实现带筛选的分组统计——聚类DataFrame分析需求
问题:PySpark按聚类分组生成变量统计指标
原始数据结构
现有聚类后的PySpark DataFrame结构如下:
| Cluster | Variable 1 | Variable 2 |
|---|---|---|
| 0 | 334 | 32 |
| 0 | 0 | 45 |
| 3 | 453 | 0 |
| 3 | 320 | 0 |
| 0 | 0 | 28 |
| 1 | 467 | 49 |
| 3 | 324 | 16 |
| 1 | 58 | 2 |
需求说明
需要针对每个Cluster和每个Variable生成以下统计字段:
%of0:该Cluster内当前变量值为0的记录占比(百分比,取整)%ofvals != 0:该Cluster内当前变量值不为0的记录占比(百分比,取整)Count of vals != 0:该Cluster内当前变量值不为0的记录数Sum of values:该Cluster内当前变量的总和%universe:该Cluster内当前变量总和占全局该变量总和的百分比(取整)
示例统计结果
Variable 1统计结果
| Cluster | %of0 | %ofvals != 0 | Count of vals != 0 | Sum of values | %universe |
|---|---|---|---|---|---|
| 0 | 67 | 33 | 1 | 334 | 17 |
| 1 | 0 | 100 | 2 | 525 | 27 |
| 3 | 0 | 100 | 3 | 1097 | 56 |
Variable 2统计结果
| Cluster | %of0 | %ofvals != 0 | Count of vals != 0 | Sum of values | %universe |
|---|---|---|---|---|---|
| 0 | 0 | 100 | 3 | 105 | 61 |
| 1 | 0 | 100 | 2 | 51 | 29 |
| 3 | 67 | 33 | 1 | 16 | 10 |
注:%universe计算示例:Variable 1全局总和为334+0+453+320+0+467+324+58=1956,Cluster 0的总和334占比为(334/1956)*100≈17%。
尝试的代码(存在逻辑问题)
for i in list_of_variables: print(i) df.groupBy('Cluster').agg((count((col(i) == 0) / df.filter(col('Cluster') == 0).count()) * 100).alias('% of 0'), (count((col(i) != 0) / df.filter(col('Cluster') == 0).count() * 100).alias('% of vals diff than 0')..
解决方案
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化SparkSession(若已初始化可跳过) spark = SparkSession.builder.appName("ClusterStats").getOrCreate() # 加载示例数据(实际使用时替换为你的DataFrame) data = [ (0, 334, 32), (0, 0, 45), (3, 453, 0), (3, 320, 0), (0, 0, 28), (1, 467, 49), (3, 324, 16), (1, 58, 2) ] df = spark.createDataFrame(data, ["Cluster", "Variable 1", "Variable 2"]) # 定义需要统计的变量列表 list_of_variables = ["Variable 1", "Variable 2"] # 预计算每个变量的全局总和,用于后续%universe计算 global_sums = {} for var in list_of_variables: total_sum = df.agg(F.sum(var)).collect()[0][0] global_sums[var] = total_sum # 循环处理每个变量,生成统计结果 for var in list_of_variables: print(f"\n=== {var} 统计结果 ===") # 按Cluster分组计算核心统计指标 cluster_stats = df.groupBy("Cluster").agg( # 计算%of0:组内0值记录占比 F.round(F.sum(F.when(F.col(var) == 0, 1).otherwise(0)) / F.count("*") * 100).alias("%of0"), # 计算%ofvals !=0:组内非0值记录占比 F.round(F.sum(F.when(F.col(var) != 0, 1).otherwise(0)) / F.count("*") * 100).alias("%ofvals != 0"), # 计算非0值记录数 F.sum(F.when(F.col(var) != 0, 1).otherwise(0)).alias("Count of vals != 0"), # 计算组内变量总和 F.sum(var).alias("Sum of values") ) # 计算%universe:组内总和占全局总和的比例 cluster_stats = cluster_stats.withColumn( "%universe", F.round(F.col("Sum of values") / global_sums[var] * 100) ) # 按Cluster排序后展示结果 cluster_stats.orderBy("Cluster").show()
代码逻辑说明
- 全局总和预计算:提前算出每个变量的全局总和,避免重复计算,提升效率
- 分组统计核心指标:
- 用
F.when标记0/非0记录,求和得到对应数量 - 用数量除以组内总记录数,乘以100后取整得到占比
- 直接用
F.sum计算组内变量总和
- 用
- 全局占比计算:将组内总和除以预计算的全局总和,取整得到
%universe - 循环处理变量:遍历所有目标变量,输出每个变量的聚类统计结果
内容的提问来源于stack exchange,提问作者Loki
相关产品推荐
相关产品推荐

