【发布时间】:2021-11-17 12:01:33
【问题描述】:
我需要用新数据对旧数据执行更新插入(Upsert)。
伪代码:
old_data = spark.read.parquet('s3://bucket/old_data/')
new_data= spark.read.parquet('s3://bucket/new_data/')
common_records = old_data.join(new_data,on=opk,how="inner")
non_match_records = old_data.join(new_data,on=opk,how="left_anti")
new_records = new_data.join(old_data,on=opk,how="left_anti")
dfs = [common_records , non_match_records , new_records ]
final_data = reduce(DataFrame.unionAll, dfs)
final_data .cache()
final_data.write.parquet('s3://bucket/old_data/')
错误:
即使我缓存了数据,还在寻找old_data路径,有没有办法直接写入旧数据s3路径。
我已经尝试将它写入一些临时路径并从中读取并写入主路径,如下所示它有效,但是当我在 Billons 中有数据时,它需要时间处理。
final_data.write.parquet('s3://bucket/temp/')
df = spark.read.parquet('s3://bucket/temp/')
df.write.parquet('s3://bucket/old_data/')
我想减少这个临时读写部分。
提前致谢:)
【问题讨论】:
-
如果只是将旧数据和新数据合并,为什么不使用
append写模式直接在旧数据目录中写入新数据呢?new_data.write.parquet('s3://bucket/old_data/', 'append') -
我们不能这样做@VincentDoba 我已经用实际代码更新了问题。检查一次。我们不能将保存模式用作“追加”
标签: apache-spark amazon-s3 pyspark apache-spark-sql