You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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,核心原因有两个:

  1. Parquet是列式存储格式,天生支持只读取指定列,不需要加载整个文件。
  2. Spark的Catalyst优化器会做**列裁剪(Column Pruning)**优化:当你执行select column_1, column_4, column 10 from table_name时,优化器会把“只读取这几列”的逻辑下推到Parquet数据源读取阶段,直接跳过其他列的读取。

如果你想验证这一点,可以在代码里加上df.explain(true),查看输出的物理执行计划,你会看到Parquet扫描的部分明确列出了这几个目标列。

问题3:Parquet vs JDBC(MySQL)场景的执行流程差异,以及为什么JDBC的read阶段耗时更长?

这两种场景的核心差异在于数据源的特性和Spark的优化策略不同,咱们分别拆解流程:

Parquet场景的完整执行流程

  1. 逻辑计划定义阶段(无实际计算):
    • spark.read.parquet(path):生成逻辑计划,记录数据来源是指定路径的Parquet文件,无实际IO。
    • createOrReplaceTempView("table_name"):只是把这个逻辑计划注册成临时视图,依然没有实际操作。
    • spark.sql("select ..."):基于之前的逻辑计划生成新的逻辑计划,加入了列裁剪的规则,但还是不执行。
  2. 动作触发阶段(真正执行):
    • 当你调用println(df.count())时,Spark的Catalyst优化器会把逻辑计划转换成最优的物理计划,把列裁剪逻辑下推到Parquet读取层。
    • 分布式读取Parquet文件的指定列,计算行数后返回结果——这一步才是真正的IO和计算阶段。

JDBC(MySQL)场景的完整执行流程

你当前的代码写法会导致spark.read阶段就完成了全量数据拉取,这是JDBC数据源的特性导致的:

  1. spark.read.jdbc阶段(立即执行全量读取):
    • 如果你传入的query是类似"select * from mysql_table"的全表查询,Spark会立刻执行这个JDBC查询,把MySQL中的所有数据拉取到Spark的Executor节点上,这一步是同步执行的,不属于懒加载。
    • 所以你会看到spark.read阶段耗时很长——因为这一步已经在做跨网络的全量数据抽取了。
  2. 后续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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 04:20:33