Spark任务节点无法访问主目录时如何导入自定义Python模块
问题根源
Spark Driver端的代码运行在提交任务的节点,你仅在Driver的系统路径中添加了本地/home/hadoop/目录,也只在Driver节点本地存在module.py文件;而UDF逻辑会分发到各个Executor节点的进程中运行,Executor既没有/home/hadoop/目录的访问权限,也没有收到module.py文件,因此会抛出模块不存在的错误。
解决方案
共有三种常用的解决方案,按需选择即可:
方案1:提交任务时通过
--py-files参数传递模块
提交Spark任务时新增--py-files参数指定自定义模块路径,Spark会自动将模块分发到所有Executor的工作目录,Executor可直接导入使用,无需额外配置系统路径。
提交命令示例:spark-submit --py-files /home/hadoop/module.py /home/hadoop/main.py代码调整:删除原代码中
sys.path.append('/home/hadoop/')这一行即可。方案2:代码内动态分发模块
如果不想修改任务提交命令,可以在初始化SparkSession后,调用addPyFile方法主动将模块分发到所有Executor,效果和方案1完全一致。
只需在原代码中新增一行配置即可,无需修改提交命令:spark.sparkContext.addPyFile('/home/hadoop/module.py')方案3:多模块场景打包分发
如果你有多个自定义Python模块需要导入,可以将所有模块打包为.egg或者.whl格式的安装包,再通过上述--py-files参数或者addPyFile方法传入即可。
修改后可正常运行的完整代码示例
import pyspark.sql.functions as F import pyspark.sql.types as T from pyspark.sql import SparkSession spark = SparkSession.builder.enableHiveSupport().getOrCreate() # 使用方案2时保留此行,使用方案1可删除此行 spark.sparkContext.addPyFile('/home/hadoop/module.py') import module if __name__ == '__main__': df = spark.createDataFrame([['a', 1], ['b', 2]], schema=['id', 'value']) df.show() print(module.incr(5)) incr_udf = F.udf(lambda val: module.incr(val), T.IntegerType()) df = df.withColumn('new_value', incr_udf('value')) df.show()
内容的提问来源于stack exchange,提问作者Thirupathi Thangavel
相关产品推荐
相关产品推荐

