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

执行器未装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的逻辑根本没有分发到执行器节点运行,具体拆分如下:

  1. 分发到dn节点(执行器)运行的逻辑,完全不依赖pandas
    代码里textFile、filter、map都是RDD转换算子,逻辑序列化后发到执行器上运行,这部分逻辑只用到了Python内置的字符串split方法,没有任何pandas相关调用,执行器不需要安装pandas就能跑完这部分计算。
  2. countByValue()的结果最终会聚合回Driver端
    作为行动算子,countByValue()会先在各执行器本地完成分片内的日期值计数,再把所有分片的计数结果拉回到提交任务的Driver进程(也就是运行spark-submit的网关节点),最终返回的是一个普通Python字典,整个过程不会在执行器侧调用pandas。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 03:12:34