Spark DataFrame何时从数据源拉取数据?不同查询collect耗时差异大
核心原理说明
首先明确Spark的核心执行机制:所有生成DataFrame的操作都属于懒执行的转换操作,spark_session.sql()执行时仅完成SQL语法校验、逻辑执行计划生成、元数据加载的工作,不会实际从数据源拉取任何数据,也不会执行查询逻辑。只有调用collect()这类行动操作时,才会触发整个查询逻辑的实际执行,因此collect()的耗时完全对应整个SQL查询的真实执行耗时,你观察到的spark_session.sql()固定耗时4秒就是元数据处理的固定开销,和查询本身复杂度无关。
两个查询的耗时差异原因
- 无排序的
limit 1查询
该查询不需要全局结果校验,Spark调度时仅需要随机读取任意一个数据分区的第一条数据即可返回结果,不需要扫描全表,也不需要跨节点做数据交换,整体计算量极小,因此collect()耗时仅为数十秒。 - 带
ORDER BY ingest_date的limit 1查询
该查询要求返回全局按ingest_date排序后的第一条数据,Spark必须完成以下操作才能拿到正确结果:- 扫描全表所有分区的所有数据,每个分区各自计算当前分区内
ingest_date最小的一条记录 - 将所有分区的计算结果通过网络shuffle传输到同一个Executor节点做全局排序,最终拿到全局最小的那条记录
如果items表数据量极大,全表扫描+跨节点数据传输的开销会非常高,因此最终耗时达到20分钟。
- 扫描全表所有分区的所有数据,每个分区各自计算当前分区内
优化建议
如果需要频繁执行这类取时间最早/最晚的查询,可以做以下优化:
- 给
ingest_date字段建立索引 - 对
items表按ingest_date做分区存储,查询时Spark可以直接定位到最早的分区,不需要扫描全表
内容的提问来源于stack exchange,提问作者Zeeshan Shamsuddeen
相关产品推荐
相关产品推荐

