基于年月匹配删除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
相关产品推荐
相关产品推荐

