PySpark循环使用Union合并DataFrame仅保留最后一次数据的问题
问题原因分析
你的循环逻辑存在问题:每次循环都直接用原始的sdf1和当前语言的子DF执行union,再将结果赋值给sdf_final。这意味着第一次循环后sdf_final是sdf1 + pol的行,但第二次循环时,你又用sdf1和tur的子DF合并,直接覆盖了之前的sdf_final,最终只保留了最后一次循环的结果(sdf1 + tur)。
两种解决方法
方法一:修正循环逻辑
初始化sdf_final为原始的sdf1,之后每次循环都基于当前的sdf_final去合并新的子DF,而非每次都从sdf1重新开始:
from pyspark.sql import functions as F langs = ['pol', 'tur'] # 初始化sdf_final为原始的sdf1 sdf_final = sdf1 for lang in langs: sdf_l = sdf2.where(F.col('lang') == lang) # 基于当前sdf_final合并子DF,再赋值回sdf_final sdf_final = sdf_final.union(sdf_l)
方法二:直接过滤后一次性合并(更高效)
无需循环,直接用isin方法过滤出sdf2中指定语言的所有行,再和sdf1做一次union,这种方式更简洁,还能避免多次union的性能开销:
from pyspark.sql import functions as F langs = ['pol', 'tur'] # 一次性过滤出所有指定语言的行 sdf_filtered = sdf2.where(F.col('lang').isin(langs)) # 一次union完成合并 sdf_final = sdf1.union(sdf_filtered)
如果你的Spark版本在2.3及以上,推荐使用unionByName替代union,它可以自动匹配列名顺序,避免因列顺序不一致导致的错误:
sdf_final = sdf1.unionByName(sdf_filtered)
内容的提问来源于stack exchange,提问作者Rory
相关产品推荐
相关产品推荐

