Spark DAG执行顺序问题:多Action操作的执行顺序是否固定?
Spark多Action执行逻辑说明
先给明确结论
默认场景下,Spark会严格按照你代码的书写顺序同步执行三个写入操作,也就是先执行完df1的写入,再执行df2的写入,最后执行df3的写入,不会自动并行执行无依赖的df1和df3写入任务。
原因解释
- Spark所有action操作(包括写入、count、collect等)都是阻塞型操作:driver端执行到
df1.write("")这一行时,会先向集群提交对应Job,全程等待这个Job运行结束、写入完成后,才会执行下一行代码的df2.write(""),依此类推。 - 我们常说的DAG优化、Stage并行调度,都是单个Job内部的逻辑,Spark默认不会自动识别多个独立Job之间的依赖关系,也不会主动跨Job做并行调度优化。
怎么实现无依赖任务并行执行
如果你希望df1和df3的写入并行运行,可以自己在driver端通过多线程(比如Scala的Future、Java的线程池)同时提交多个Job,此时Spark调度器会根据你配置的调度模式(默认FIFO,也可以改FAIR)分配集群资源并行执行这些任务。
额外注意点
你给出的代码里如果没有给df1显式调用cache()/persist(),那么执行df2写入的Job时,会重新从s3://onepath/读取一次源数据,不会复用df1写入Job里已经读取过的数据。
内容的提问来源于stack exchange,提问作者Shad Khan
相关产品推荐
相关产品推荐

