PySpark如何将Source列去重值转为独立列生成Yes/No格式宽表
PySpark 大数据量宽表转换高效实现方案
核心思路
基于Spark原生pivot算子实现行转列,你已经提前提取了Source去重值列表,刚好可以避免Spark默认执行的二次全表扫描统计唯一值逻辑,整体全程分布式执行,无Driver端数据拉取操作,完美适配1亿+行数据量级。
实现代码
from pyspark.sql.functions import lit # 1. 为每行数据添加固定存在标记 dedupeTable = dedupeTable.withColumn("exist_flag", lit("Yes")) # 2. 按ID分组执行行转列,直接传入已提取的去重Source列,减少一次全表扫描Shuffle result_df = dedupeTable.groupBy("ID") \ .pivot("Source", dedupeTableColumnNamesCleaned) \ .agg({"exist_flag": "first"}) # 3. 所有空值统一替换为No,对应ID无对应Source的场景 result_df = result_df.na.fill("No", subset = dedupeTableColumnNamesCleaned)
优化说明
- 若原始数据中存在同一个ID+Source重复的情况,可将聚合逻辑从
first替换为max,无需提前对原始数据做去重,进一步减少计算步骤 - 若Source去重值数量超过1000,宽表列数过多时建议调整为行式存储结果,避免宽表查询性能问题
内容的提问来源于stack exchange,提问作者Ali Shah
相关产品推荐
相关产品推荐

