【问题标题】:Spark Dataframe Filter OptimizationSpark Dataframe 过滤器优化
【发布时间】:2019-09-05 07:30:58
【问题描述】:

我正在从 s3 存储桶中读取大量文件。

读完这些文件后,我想对数据框进行过滤操作。

但是当执行过滤操作时,数据会再次从 s3 存储桶中下载。如何避免数据框重新加载?

我已尝试在过滤操作之前缓存和/或持久化数据帧。但是,仍然以某种方式再次从 s3 存储桶中提取数据。

var df = spark.read.json("path_to_s3_bucket/*.json")

df.persist(StorageLevel.MEMORY_AND_DISK_SER_2)

df = df.filter("filter condition").sort(col("columnName").asc)

如果数据帧被缓存,则不应再次从 s3 重新加载。

【问题讨论】:

  • 你能告诉我们explain的计划吗?你怎么确定这又是从存储桶中读取的?
  • 如果担心重新阅读,请对本地进行 distcp 并从本地本身读取

标签: scala apache-spark apache-spark-sql


【解决方案1】:

当你打电话时

var df = spark.read.json("path_to_s3_bucket/*.json")

幕后发生的事情是 spark 进行分区发现、文件列表和架构推断(如果您有很多文件,这可能会在后台运行 sum 作业以并行执行文件列表)。

下次调用时

df.persist(StorageLevel.MEMORY_AND_DISK_SER_2)

只有您想要持久化数据的信息被传递给查询计划,但此时没有持久化(这是一个惰性操作)。

下次调用时

df = df.filter("filter condition").sort(col("columnName").asc)

再次只更新查询计划。

现在如果调用show()count()等动作,就会处理查询计划并执行spark job。所以现在数据将被加载到集群上,它会被写入内存(因为缓存),然后从缓存中读回,根据您的查询计划对其进行过滤、排序和进一步处理。

【讨论】:

  • 您好,感谢您的意见。但是当执行 var df = spark.read.json("path_to_s3_bucket/*.json") 时,加载数据帧需要将近两个小时。然后再次进行过滤操作时,大约需要两个小时。所以我的猜测是数据框再次重新加载
  • @ChinmayR 您正在阅读多少文件?前两个小时可能实际上只是文件列表和模式推断。尝试为 DataframeReader 提供一个模式,看看它是否加快了速度。
  • 总共大约 35 个文件,总大小约为 2 GB。我在它之后使用 printSchema 和 show 命令。这给了我样本数据。并在过滤命令上再次读取数据。
  • @ChinmayR 数据不应该被过滤命令读取,操作是惰性的,它应该只更新查询计划。我不太明白这怎么可能。另外 35 个文件,总共 2GB 也不算多,读完应该不会花 2 个小时。你的集群有多大?
  • 可能是我的网速太慢了。当我尝试在 ec2 实例上运行它时,几乎不需要 5 分钟
猜你喜欢
  • 2023-03-26
  • 2019-01-07
  • 2019-10-05
  • 2022-01-23
  • 1970-01-01
  • 2020-11-29
  • 2020-06-16
  • 2021-12-09
  • 2019-02-10
相关资源
最近更新 更多