Spark DataFrame是否懒加载Parquet?Parquet与JDBC列读取及流程解析
咱们一个个来拆解你的问题,把这些Spark的执行逻辑讲清楚:
问题1:Spark DataFrame是否懒加载Parquet数据?
没错,Spark的DataFrame(包括从Parquet读取的)完全遵循**懒加载(Lazy Evaluation)**机制。简单来说:
- 当你执行
spark.read.parquet(path)时,Spark并不会立刻去磁盘读取数据,它只是创建了一个描述「数据来源+后续可能的转换逻辑」的逻辑执行计划,相当于打了个“待办清单”。 - 只有当你触发**动作(Action)**操作(比如
count()、show()、collect()这些)的时候,Spark才会把逻辑计划转换成可执行的物理计划,真正去磁盘读取数据并完成计算。
问题2:你的Parquet处理代码是否仅读取所需列?
肯定是的!这段代码会只读取SQL中指定的column_1、column_4和column 10,核心原因有两个:
- Parquet是列式存储格式,天生支持只读取指定列,不需要加载整个文件。
- Spark的Catalyst优化器会做**列裁剪(Column Pruning)**优化:当你执行
select column_1, column_4,column 10from table_name时,优化器会把“只读取这几列”的逻辑下推到Parquet数据源读取阶段,直接跳过其他列的读取。
如果你想验证这一点,可以在代码里加上df.explain(true),查看输出的物理执行计划,你会看到Parquet扫描的部分明确列出了这几个目标列。
问题3:Parquet vs JDBC(MySQL)场景的执行流程差异,以及为什么JDBC的read阶段耗时更长?
这两种场景的核心差异在于数据源的特性和Spark的优化策略不同,咱们分别拆解流程:
Parquet场景的完整执行流程
- 逻辑计划定义阶段(无实际计算):
spark.read.parquet(path):生成逻辑计划,记录数据来源是指定路径的Parquet文件,无实际IO。createOrReplaceTempView("table_name"):只是把这个逻辑计划注册成临时视图,依然没有实际操作。spark.sql("select ..."):基于之前的逻辑计划生成新的逻辑计划,加入了列裁剪的规则,但还是不执行。
- 动作触发阶段(真正执行):
- 当你调用
println(df.count())时,Spark的Catalyst优化器会把逻辑计划转换成最优的物理计划,把列裁剪逻辑下推到Parquet读取层。 - 分布式读取Parquet文件的指定列,计算行数后返回结果——这一步才是真正的IO和计算阶段。
- 当你调用
JDBC(MySQL)场景的完整执行流程
你当前的代码写法会导致spark.read阶段就完成了全量数据拉取,这是JDBC数据源的特性导致的:
spark.read.jdbc阶段(立即执行全量读取):- 如果你传入的
query是类似"select * from mysql_table"的全表查询,Spark会立刻执行这个JDBC查询,把MySQL中的所有数据拉取到Spark的Executor节点上,这一步是同步执行的,不属于懒加载。 - 所以你会看到
spark.read阶段耗时很长——因为这一步已经在做跨网络的全量数据抽取了。
- 如果你传入的
- 后续SQL和动作阶段(本地计算):
- 当你创建临时视图、执行
select和count()时,Spark只是在已经拉取到本地的数据上做列裁剪和计数,这一步的计算量远小于之前的全量数据拉取,所以耗时更短。
- 当你创建临时视图、执行
优化JDBC场景的小技巧
如果想让JDBC也实现类似Parquet的懒加载和列裁剪优化,把列选择逻辑放到JDBC的query参数里就行:
// 把需要的列写到JDBC查询里,推送给MySQL执行 val query = "(select column_1, column_4, `column 10` from mysql_table) as temp_table" spark.read.format("jdbc").jdbc(jdbcUrl, query, props).createOrReplaceTempView("table_name") spark.sql("select column_1, column_4, `column 10` from table_name") df.show() println(df.count())
这样MySQL只会返回你需要的列,spark.read阶段的耗时会大幅降低,而且整个流程也会遵循懒加载——直到触发show()或count()才会执行MySQL查询并拉取数据。
内容的提问来源于stack exchange,提问作者Devas
相关产品推荐
相关产品推荐

