【问题标题】:UPSERT in parquet Pyspark拼花地板 Pyspark 中的 UPSERT
【发布时间】:2020-05-12 07:18:03
【问题描述】:

我在 s3 中有 parquet 文件,具有以下分区: 年/月/日/ some_id 使用 Spark (PySpark),我每天都想在最后 14 天 进行 UPSERT - 我想替换 s3 中的现有数据(每个分区一个 parquet 文件),但不删除14天之前的日子.. 我尝试了两种保存模式: append - 不好,因为它只是添加了另一个文件。 overwrite - 正在删除过去的数据和其他分区的数据。

有什么方法或最佳实践可以克服这个问题吗?我应该在每次运行中从 s3 读取所有数据,然后再写回来吗?也许重命名文件以便 append 将替换 s3 中的当前文件?

非常感谢!

【问题讨论】:

    标签: amazon-s3 pyspark etl parquet


    【解决方案1】:

    据我所知,S3 没有更新操作。一旦将对象添加到 s3 就无法修改。 (要么你必须替换另一个对象或附加一个文件)

    不管您是否担心必须读取所有数据,您都可以指定要读取的时间线,分区修剪有助于仅读取时间线内的分区。

    【讨论】:

      【解决方案2】:

      我通常会做类似的事情。就我而言,我执行 ETL 并将一天的数据附加到 parquet 文件中:

      关键是使用您要写入的数据(在我的情况下是实际日期),确保按date 列分区并覆盖当前日期的所有数据。

      这将保留所有旧数据。举个例子:

      (
          sdf
          .write
          .format("parquet")
          .mode("overwrite")
          .partitionBy("date")
          .option("replaceWhere", "2020-01-27")
          .save(uri)
      )
      

      您还可以查看 delta.io,它是 parquet 格式的扩展,提供了一些有趣的功能,例如 ACID 事务。

      【讨论】:

        【解决方案3】:

        感谢大家提供有用的解决方案。 我最终使用了一些服务于我的用例的配置 - 在我编写 parquet 时使用 overwrite 模式,以及以下配置:

        我添加了这个配置:

        spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
        

        使用此配置,spark 只会覆盖它要写入数据的分区。所有其他(过去的)分区保持不变 - 请参见此处:

        https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-dynamic-partition-inserts.html

        【讨论】:

          猜你喜欢
          • 2021-04-10
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2021-03-22
          • 2018-06-19
          • 2017-12-30
          相关资源
          最近更新 更多