如何在PySpark 2.0.1中将关联重复值整合到单行生成新列?
可以在PySpark 2.0.1中实现该需求!
核心思路是利用**透视(pivot)**操作将不同月份对应的聚合结果从行转换为列,再与你的交叉表关联,最终得到每个name一行的目标结构。下面是基于你现有代码的完整实现步骤:
步骤1:准备透视所需的月份列表
先提取所有唯一的月份值,避免pivot时全表扫描(PySpark 2.0.1支持指定pivot的values参数,提升性能):
# 获取所有唯一月份 months = degree_df.select("month").distinct().rdd.flatMap(lambda x: x).collect()
步骤2:对聚合表进行透视转换
针对你已有的table_count_d,按name分组,将month作为透视列,同时保留每个月份对应的month值、min(degree)和max(degree):
# 透视操作:将行转列,每个月份生成一组专属列 pivoted_degree = table_count_d.groupBy("name") \ .pivot("month", months) \ .agg( first("month").alias("month"), # 提取对应月份名 min("degree").alias("min_degree"), max("degree").alias("max_degree") )
透视后,列名会自动变为[月份]_[字段名]的格式(比如May_month、April_min_degree),保证列名唯一(PySpark不允许重复列名,这是合理的规范)。
步骤3:关联交叉表并处理空值
将透视后的表与你的交叉表table_count_c左连接,再填充空值(比如Emma的April相关数据为空,需要补0):
# 左连接保留所有name table_final = table_count_c.join(pivoted_degree, on="name", how="left_outer") # 填充min/max的空值为0.0 table_final = table_final.fillna(0.0, subset=[col for col in table_final.columns if "min_degree" in col or "max_degree" in col]) # 填充空的month列为对应月份名(比如Emma的April_month为空,填充"April") for month in months: table_final = table_final.withColumn( f"{month}_month", when(col(f"{month}_month").isNull(), month).otherwise(col(f"{month}_month")) )
步骤4:调整列顺序匹配期望结构
指定列顺序来对齐你想要的展示效果:
desired_columns = [ "name", "April", "May", "May_month", "May_min_degree", "May_max_degree", "April_month", "April_min_degree", "April_max_degree" ] table_final = table_final.select(desired_columns)
最终结果展示
执行table_final.show()后会得到:
+-----+-----+---+----------+---------------+---------------+------------+-----------------+-----------------+ | name|April|May|May_month|May_min_degree|May_max_degree|April_month|April_min_degree|April_max_degree| +-----+-----+---+----------+---------------+---------------+------------+-----------------+-----------------+ |Ahmad| 2| 1| May| 38.0| 38.0| April| 40.0| 49.0| | Emma| 0| 2| May| 45.0| 50.0| April| 0.0| 0.0| +-----+-----+---+----------+---------------+---------------+------------+-----------------+-----------------+
关键说明
- PySpark 2.0.1完全支持
pivot操作(该特性在PySpark 1.6版本就已引入),所以这个方案是可行的。 - 透视操作是实现“将重复值对应数据合并到同一行”的核心,它把多行的月份数据转换为单行的多列结构。
- 我们保留了唯一列名(而非你期望的重复列名),因为PySpark的DataFrame不允许重复列名,这是数据库的通用规范,也更利于后续数据处理。
内容的提问来源于stack exchange,提问作者Ahmad Senousi
相关产品推荐
相关产品推荐

