执行器未装Python模块时spark-submit为何仍能成功运行PySpark任务?
问题场景
- 集群共7个节点,角色分配如下:
- 1: master
- 2: nn1
- 3: nn2
- 4: dn1
- 5: dn2
- 6: dn3
- 7: gn(网关节点)
- 环境前提:仅在网关节点安装了任务所需的Python依赖,通过
spark-submit提交PySpark任务,预期任务会因为dn节点(执行器节点)缺少Python模块失败,但实际任务运行成功。 - 提交的任务代码如下:
from pyspark import SparkConf, SparkContext import pandas as pd # Spark configuration conf = SparkConf().setAppName("uber-date-trips") sc = SparkContext(conf=conf) # data parsing lines = sc.textFile( "/user/tester/trips_2020-03.csv" ) header = lines.first() filtered_lines = lines.filter(lambda row:row != header) # data filtering and counting same date data dates = filtered_lines.map(lambda x: x.split(",")[2].split(" ")[0]) result = dates.countByValue() # save the result to csv type pd.Series(result, name="trips").to_csv("trips_date.csv")
任务未失败的核心原因
本质是对PySpark的代码执行位置判断有误,用到pandas的逻辑根本没有分发到执行器节点运行,具体拆分如下:
- 分发到dn节点(执行器)运行的逻辑,完全不依赖pandas
代码里textFile、filter、map都是RDD转换算子,逻辑序列化后发到执行器上运行,这部分逻辑只用到了Python内置的字符串split方法,没有任何pandas相关调用,执行器不需要安装pandas就能跑完这部分计算。 countByValue()的结果最终会聚合回Driver端
作为行动算子,countByValue()会先在各执行器本地完成分片内的日期值计数,再把所有分片的计数结果拉回到提交任务的Driver进程(也就是运行spark-submit的网关节点),最终返回的是一个普通Python字典,整个过程不会在执行器侧调用pandas。- pandas相关逻辑完全在Driver端执行
最后pd.Series(result, name="trips").to_csv("trips_date.csv")这行代码,是拿到聚合完成的结果字典后,在网关节点的Driver进程里本地运行的,只需要网关节点装了pandas就能正常执行,和执行器节点有没有装pandas没有任何关系。
补充说明:只有当你在map、filter、UDF这类需要分发到执行器运行的逻辑里调用了pandas,或者使用了pandas UDF特性时,才要求所有执行器节点都安装对应版本的pandas依赖,否则仅Driver端安装即可满足运行要求。
内容的提问来源于stack exchange,提问作者JiHooney
相关产品推荐
相关产品推荐

