spark.read与spark.sql是否为懒转换?源数据更新后调用行动算子为何返回旧结果?
问题1:spark.read和spark.sql是否属于懒转换操作?
二者均属于懒转换操作:调用spark.read系列读取接口或者spark.sql()生成DataFrame时,Spark只会记录对应的读取逻辑、SQL逻辑,不会立刻触发实际的IO和计算,只有对生成的DataFrame调用行动算子时,才会真正执行对应逻辑。
问题2:场景中的认知偏差说明
你的基础认知「触发行动算子时Spark会基于DAG执行全链路操作,包含读取步骤」本身没有错误,你遇到的返回旧数据的问题,是忽略了Spark的表元数据、文件列表快照机制:
- 你执行
df = spark.sql("select * from dummy.table1")生成DataFrame时,Spark为了优化后续查询性能,会提前解析该表的元数据,扫描该表对应的所有数据文件,把这一时刻的文件列表作为快照写入DataFrame的逻辑计划中,后续所有对这个DataFrame的操作,都是基于这份固定的文件列表执行的。 - 你后续插入的新记录会对应生成新的数据文件,但旧DataFrame的逻辑计划里没有包含新文件路径,所以第二次调用
df.count()时,仍然只会读取创建DataFrame时扫描到的旧文件,自然返回旧结果2,而非你预期的3。
认知偏差核心点
你默认Spark每次执行DAG都会重新解析外部表的最新元数据、拉取最新的文件列表,但实际上DataFrame一旦创建,其对应的逻辑计划就已经固定,不会在每次行动算子触发时动态更新外部依赖的元数据和文件列表。
要实现你预期的动态读取最新数据的效果,可选择以下两种方案:
- 每次需要查询最新数据前,重新执行一次
df = spark.sql("select * from dummy.table1"),生成包含最新文件列表快照的新DataFrame - 查询前先执行
spark.sql("REFRESH TABLE dummy.table1")刷新表的元数据缓存,再重新生成DataFrame
内容的提问来源于stack exchange,提问作者ROHIT CHAUHAN
相关产品推荐
相关产品推荐

