You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.01 04:15:01