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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:55:24