如何在PySpark中按ID汇总DataFrame并取各列首个非空值?
PySpark实现按ID汇总取各列第一个非空值
原始DataFrame结构如下(实际包含更多行和列):
| id | col1 | col2 | col3 |
|---|---|---|---|
| 1 | A | null | null |
| 1 | null | B | null |
| 1 | null | null | C |
| 2 | D | null | null |
| 2 | null | E | null |
| 2 | null | null | F |
需求:按ID汇总DataFrame,为每个ID的所有列选取第一个非空值。
期望结果:
| id | col1 | col2 | col3 |
|---|---|---|---|
| 1 | A | B | C |
| 2 | D | E | F |
在R的dplyr中可通过以下代码实现:
df %>% group_by(id) %>% summarize_all(list(~first(na.omit(.))))
PySpark中的实现方法
方法1:使用first函数(推荐,高效)
PySpark的first函数支持ignoreNulls=True参数,可直接忽略空值,取分组内该列的第一个非空值。
如果列数较少,可手动指定每个列的聚合逻辑:
from pyspark.sql import functions as F # 假设你的DataFrame名为df result_df = df.groupBy("id").agg( F.first("col1", ignoreNulls=True).alias("col1"), F.first("col2", ignoreNulls=True).alias("col2"), F.first("col3", ignoreNulls=True).alias("col3") ) result_df.show()
如果列数较多,可动态生成聚合表达式,避免重复代码:
from pyspark.sql import functions as F # 获取除id外的所有列 target_cols = [col for col in df.columns if col != "id"] # 生成每个列的first聚合表达式 agg_expressions = [F.first(col, ignoreNulls=True).alias(col) for col in target_cols] # 执行分组聚合 result_df = df.groupBy("id").agg(*agg_expressions) result_df.show()
方法2:使用coalesce+collect_list
这种方法会先收集分组内该列的所有值,再通过coalesce取出第一个非空值,效率低于方法1(因为需要收集所有值),适合特殊场景:
from pyspark.sql import functions as F target_cols = [col for col in df.columns if col != "id"] agg_expressions = [F.coalesce(*F.collect_list(col)).alias(col) for col in target_cols] result_df = df.groupBy("id").agg(*agg_expressions) result_df.show()
内容的提问来源于stack exchange,提问作者Megha
相关产品推荐
相关产品推荐

