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

基于年月匹配删除DataFrame旧数据并插入新数据的PySpark问题

PySpark实现数据更新:匹配替换+新增未匹配数据

核心思路

要实现「保留file1未匹配数据、替换匹配数据、新增file2未匹配数据」的需求,最简洁的方式是通过左反连接筛选出file1中未被file2匹配的记录,再将其与file2的全部数据合并。具体逻辑:

  • 左反连接(left_anti):提取file1中不存在于file2(按year、month匹配)的行
  • 合并(union):将上述筛选结果与file2的所有行合并,即可得到最终的更新后数据

完整代码实现

from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder.appName("DataUpdate").getOrCreate()

# 加载file1和file2数据到DataFrame(替换为你的实际文件路径)
df_file1 = spark.read.csv("file1.csv", header=True, inferSchema=True)
df_file2 = spark.read.csv("file2.csv", header=True, inferSchema=True)

# 左反连接筛选file1中未被file2匹配的记录
df_file1_unmatched = df_file1.join(
    df_file2,
    on=["year", "month"],
    how="left_anti"
)

# 合并未匹配记录与file2全量数据,生成最终结果
final_df = df_file1_unmatched.union(df_file2)

# 按年份月份排序查看结果
final_df.orderBy("year", "month").show()

结果验证

针对你的示例数据,执行代码后输出如下:

+----+-----+------+
|year|month|Amount|
+----+-----+------+
|2020|March|    30|
|2021|  Jan|    20|
|2022|  Jan|   220|
|2022|  Feb|   130|
|2022|  Aug|    12|
|2022|  Oct|   100|
+----+-----+------+

完全符合期望:保留了file1中2020-March、2021-Jan、2022-Aug的未匹配数据,替换了2022-Oct的旧数据,新增了file2中的2022-Jan、2022-Feb数据。

关于你尝试的exists方式

如果想用exists实现,需要构造子查询判断当前行是否在file2中存在匹配,但代码更繁琐,还容易因作用域问题出错,示例如下(不推荐使用,不如左反连接直观):

from pyspark.sql.functions import exists, col

# 构造子查询
file2_subquery = df_file2.select("year", "month").alias("f2")

# 筛选file1中未匹配的行
df_file1_unmatched = df_file1.filter(
    ~exists(file2_subquery, (col("f2.year") == col("year")) & (col("f2.month") == col("month")))
)

# 合并得到结果
final_df = df_file1_unmatched.union(df_file2)

内容的提问来源于stack exchange,提问作者codetech

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 22:25:21