【问题标题】:If you save a DataFrame to disk, will Spark load that data if you use that DataFrame lower in the script?如果您将 DataFrame 保存到磁盘,如果您在脚本中使用该 DataFrame,Spark 会加载该数据吗?
【发布时间】:2019-07-02 18:01:48
【问题描述】:

如果您加载一些数据,计算一个 DataFrame,将其写入磁盘,然后稍后使用 DataFrame...假设它还没有缓存在 RAM 中(假设它还不够),Spark 会足够聪明吗从磁盘加载数据而不是从原始数据重新计算 DataFrame?

例如:

df1 = spark.read.parquet('data/df1.parquet')
df2 = spark.read.parquet('data/df2.parquet')

joined = df1.join(df2, df1.id == df2.id)
joined.write.parquet('data/joined.parquet')

computed = joined.select('id').withColummn('double_total', 2 * joined.total)
computed.write.parquet('data/computed.parquet')

在适当的情况下,当我们存储computed时,它会从data/joined.parquet加载joined DataFrame,还是总是通过加载/加入df1/df2重新计算,如果不是当前缓存joined

【问题讨论】:

  • 如此有效地询问 write.parquet 是否算作连接转换的操作。我也很想知道答案

标签: python apache-spark caching pyspark apache-spark-sql


【解决方案1】:

joined 数据框指向df1.join(df2, df1.id == df2.id)。据我所知,parquet writer 不会对该引用进行任何更改,因此为了加载 parquet 数据,您需要使用spark.reader.parquet(...) 构造一个新的 Spark 阅读器。

您可以从 DataFrameWriter 代码(检查 parquet/save 方法)验证上述声明,该代码返回 Unit 并且不会以某种方式修改源数据帧的引用。最后在上面的例子中回答你的问题,加入的数据框将计算一次 joined.write.parquet('data/joined.parquet') 和一次 computed.write.parquet('data/computed.parquet')

【讨论】:

  • 谢谢!这是令人惊讶的。如果它必须重新计算,我本来希望它使用磁盘上的版本,如果它认为这会更快。
  • 好吧,据我所知,使用镶木地板的唯一方法是用已经提到的阅读器重新加载它。还要考虑这样一个事实,即在写入 parquet 后,此数据存储在内存和/或最终分区中,因此,如果您不使用任何广泛的转换,Spark 将不会从头开始加载所有内容,而是会使用存储的信息和现有的沿袭
  • 是的,我意识到这一点。但是我经常测试我正在使用的机器的 RAM 的限制,所以它经常需要重新计算,并且当它这样做时从那里加载会很酷。从 Parquet 加载一个宽 DataFrame 中的几列可以使计算变得可行,否则根本不可能。保存后我会继续手动加载,谢谢!
  • 好的也可以尝试缓存它并在检查点之后立即进行。
猜你喜欢
  • 2016-01-15
  • 1970-01-01
  • 2019-11-03
  • 2015-12-23
  • 1970-01-01
  • 2017-11-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多