【问题标题】:Spark Reading and Writing to same S3 Path Giving Unable to infer Schema ErrorSpark读取和写入相同的S3路径导致无法推断架构错误
【发布时间】: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


【解决方案1】:

您需要执行一个操作来触发数据帧缓存。因此您应该修改代码 sn-p 的最后几行,如下所示:

...
final_data = final_data.cache()
final_data.count()
final_data.write.parquet('s3://bucket/old_data/')

通过对您的数据帧执行count 操作,您可以触发缓存过程,然后能够在您读取的同一目录中写入。

但是,当数据帧对于内存来说太大时,我不知道这是否会提高应用程序的性能,因为缓存回退到磁盘写入。如果您的用例是更新 parquet 文件,我建议您查看为解决此问题而创建的 DeltaLake

【讨论】:

    猜你喜欢
    • 2018-12-05
    • 1970-01-01
    • 1970-01-01
    • 2018-06-29
    • 2021-10-29
    • 1970-01-01
    • 2020-01-17
    • 2015-10-24
    • 1970-01-01
    相关资源
    最近更新 更多