如何在PySpark DataFrame中判断列全值相等并删除无信息常量列
PySpark删除全列取值相同列的最优实现方案
你原来的两种实现都存在明显的性能或逻辑问题:
- 最大值最小值对比法:每列触发2次action作业,列数多的时候执行效率极低;且存在逻辑漏洞,全为null的列min和max均为null,null等值判断不生效,会漏删全null列
- distinct计数法:每列触发1次action作业,虽然修复了全null列的删除问题,但多列场景下仍会触发多次作业,资源消耗高、速度慢
最优实现方案
最优方案是仅触发1次全局聚合作业,一次性计算所有列的唯一值计数,再批量删除符合条件的列,不管有多少列都只会跑一次Spark作业,执行效率提升幅度和列数正相关,代码也更简洁:
from pyspark.sql import functions as F # 单次聚合计算所有列的唯一值数量 col_distinct_cnt = df.agg(*[F.countDistinct(col).alias(col) for col in df.columns]).collect()[0].asDict() # 筛选需要删除的列:唯一值数为1对应全列非null且取值相同,若要同时删除全null列可将条件调整为 cnt <= 1 cols_to_drop = [col for col, cnt in col_distinct_cnt.items() if cnt == 1] # 批量删除目标列 df = df.drop(*cols_to_drop)
超大数据集优化变种
如果是TB级以上的超大规模数据集,允许极小的统计误差,可以用近似去重计数替换精确计数,进一步提升执行速度:
# 第二个参数为可接受的误差阈值,默认值为0.05,可根据业务需求调整 col_distinct_cnt = df.agg(*[F.approx_count_distinct(col, 0.01).alias(col) for col in df.columns]).collect()[0].asDict()
内容的提问来源于stack exchange,提问作者Zaka Elab
相关产品推荐
相关产品推荐

