PySpark UDF调用geodesic遇ModuleNotFoundError问题求助
问题解决:PySpark UDF中
geopy模块找不到的问题 问题背景
现有如下PySpark DataFrame:
df = spark.createDataFrame([(['target'], [2], [2], [3], [3]), (['NJ'],[3],[3], [4], [4]), (['target', 'target'],[4,5], [4,5], [6,7], [6,7]), (['CA'],[5],[5], [6], [6]), ], ('group_name', 'long', 'lat','com_long','com_lat'))
需要提取group_name列中'target'对应位置的数据,通过UDF计算地理距离。执行步骤包括获取目标索引、生成辅助列、过滤空值、定义并调用基于geopy.geodesic的UDF,但运行时出现错误:
ModuleNotFoundError: No module named 'geopy'
已确认本地安装了geopy、导入语句正确且完成空值过滤,问题仍存在。
核心原因
PySpark的Driver节点(本地运行代码的节点)和Worker节点(执行UDF等分布式任务的节点)是分离的。你在Driver端通过pip安装了geopy,但Worker节点的Python环境中没有安装该包,而UDF的逻辑是在Worker节点上执行的,因此会触发模块找不到的错误。
解决方案
方案一:提交作业时上传geopy包
将geopy打包为wheel文件,在提交Spark作业时通过--py-files参数上传到所有Worker节点:
- 先下载
geopy的wheel包:pip download geopy --no-deps -d /path/to/save - 提交作业时指定该包:
spark-submit --py-files /path/to/save/geopy-2.4.0-py3-none-any.whl your_script.py
方案二:在所有Worker节点安装geopy
直接在集群的每个Worker节点上执行安装命令:
pip install geopy
如果是大规模集群,可通过Ansible、SaltStack等自动化工具批量执行安装操作。
方案三:配置Spark使用统一的Python环境
如果集群使用自定义Python环境(如Anaconda),确保所有Worker节点的该环境中已安装geopy,并在Spark配置中指定该Python路径:
spark-submit --conf spark.pyspark.python=/path/to/your/python/environment/bin/python your_script.py
内容的提问来源于stack exchange,提问作者maketew
相关产品推荐
相关产品推荐

