Spark Driver是否执行同步阻塞调用?执行顺序与等待机制确认
Spark Driver对Action的调用是否为同步阻塞?
结论:Spark Driver对Action的调用是同步阻塞的——Driver会等待当前Action对应的Job在Executor上执行完成、拿到结果后,才会继续执行后续的代码指令。
针对你的示例解释:
- 第一个多Action示例:
spark .read .schema(mySchema) .json(myFilePath) .withColumn("a", col("b") * 2) .filter(col("c") > 300) .count() spark .read .schema(mySchema2) .json(myFilePath2) .filter(col("d") < 100) .count()
两个count()动作会严格按代码声明的顺序执行:Driver会先调度第一个count()对应的Job,等待其完全执行完毕、拿到结果后,才会触发第二个count()的调度与执行,不会并行处理这两个Job。
- 第二个带Scala语句的示例:
val df1 = spark .read .schema(mySchema) .json(myFilePath) .withColumn("a", col("b") * 2) .filter(col("c") > 300) // no execution happened until here val df1Count = df1.count() // "count" action triggers the execution println(s"df1 contains ${df1Count} rows.") // rows are logged correctly
当执行df1.count()时,Driver会触发对应的Job执行,在此期间Driver会处于阻塞状态,直到Executor完成计算并将结果返回给Driver,赋值给df1Count后,才会继续执行后续的println语句——这也是你能正确打印出结果的核心原因,因为df1Count已经拿到了计算后的数值。
补充要点:
- Spark默认的Driver进程是单线程的,所有Action的触发逻辑都在这个单线程中按顺序同步执行,不存在Spark自动异步调度Action的情况(如果需要并行执行多个Job,需要你在应用层通过多线程等方式自行实现)。
- 延迟计算(Lazy Evaluation)仅针对转换操作(Transformation),这类操作只会构建逻辑执行计划,不会立即触发计算;而Action操作是触发实际计算的入口,且每个Action的执行对Driver来说都是阻塞式的。
内容的提问来源于stack exchange,提问作者Slph
相关产品推荐
相关产品推荐

