【发布时间】:2017-07-26 19:16:09
【问题描述】:
我正在运行 emr-5.2.0,并且在 S3 中存储了一年的数据作为按天分区的 Parquet。查询一个月时,我希望 Spark 仅将一个月的数据加载到内存中。但是,我的集群内存使用情况看起来像是在加载整年 1.7TB 的数据。
我假设我可以像这样加载完整的数据湖
val lakeDF = spark.sqlContext.read.parquet("s3://mybucket/mylake.parquet")
lakeDF.cache()
lakeDF.registerTempTable("sightings")
Spark 将使用查询中的日期来仅选择与 WHERE 过滤器匹配的分区。
val leftDF = spark.sql("SELECT * FROM sightings WHERE DATE(day) BETWEEN "2016-01-09" AND "2016-01-10"")
val audienceDF = leftDF.join(ghDF, Seq("gh9"))
audienceDF.select( approxCountDistinct("device_id", red = 0.01).as("distinct"), sum("requests").as("avails") ).show()
我很好奇将分区转换为 DATE 是否会导致此问题?
我还在同一个数据集上使用 Athena/PrestoDB 运行了一些测试,很明显只扫描了几 GB 的数据。
Spark 有什么方法可以告诉我在提交查询之前要加载多少数据?
【问题讨论】:
-
您是否尝试过删除
lakeDF.cache()语句?你还冷研究了df.explain()在你的转换结束时给出的物理计划(在调用动作之前),也许这会给你一个提示 -
是的
lakeDF.cache()是问题所在。
标签: apache-spark amazon-emr parquet