【问题标题】:Bundling Python3 packages for PySpark results in missing imports为 PySpark 捆绑 Python3 包会导致缺少导入
【发布时间】:2018-07-24 00:39:09
【问题描述】:

我正在尝试运行依赖于某些 python3 库的 PySpark 作业。 我知道我可以在 Spark 集群上安装这些库,但由于我将集群重用于多个作业,我宁愿捆绑所有依赖项并通过 --py-files 指令将它们传递给每个作业。

为此,我使用:

pip3 install -r requirements.txt --target ./build/dependencies
cd ./build/dependencies
zip -qrm . ../dependencies.zip

这有效地压缩了所需包中的所有代码,以便在根级别使用。

在我的main.py 中,我可以导入依赖项

if os.path.exists('dependencies.zip'):
    sys.path.insert(0, 'dependencies.zip')

并将 .zip 添加到我的 Spark 上下文中

sc.addPyFile('dependencies.zip')

到目前为止一切顺利。

但由于某种原因,这将在 Spark 集群上产生某种依赖地狱

例如跑步

spark-submit --py-files dependencies.zip main.py

main.py(或班级)我想在哪里使用熊猫。会触发这个错误的代码:

Traceback(最近一次调用最后一次):

文件“/Users/tomlous/Development/Python/enrichers/build/main.py”,第 53 行,在 job_module = importlib.import_module('spark.jobs.%s' % args.job_name) ...

文件“”,第 978 行,在 _gcd_import 中

文件“”,第 961 行,在 _find_and_load 中

文件“”,第 950 行,在 _find_and_load_unlocked 中

文件“”,第 646 行,在 _load_unlocked 中

文件“”,第 616 行,在 _load_backward_compatible

文件“dependencies.zip/spark/jobs/classify_existence.py”,第 9 行,

文件“dependencies.zip/enrich/existence.py”,第 3 行,在

文件“dependencies.zip/pandas/init.py”,第 19 行,

ImportError:缺少必需的依赖项 ['numpy']

看着熊猫的__init__.py 我看到类似__import__(numpy)的东西

所以我假设 numpy 没有加载。

但是,如果我将代码更改为显式调用 numpy 函数,它实际上会找到 numpy,但不是它的某些依赖项

import numpy as np
a = np.array([1, 2, 3])

代码返回

Traceback(最近一次调用最后一次):

文件“dependencies.zip/numpy/core/init.py”,第 16 行,

ImportError: cannot import name 'multiarray'

所以我的问题是:

我应该如何将 python3 库与我的 spark 作业捆绑在一起,而不必在 Spark 集群上 pip3 安装所有可能的库?

【问题讨论】:

  • 我将从记录 sys.path 开始了解 pyspark 节点看到的内容。
  • 我将连接到工作节点,启动 python 并开始执行导入。
  • 我会尝试日志记录,但连接到工作节点会破坏目的。 GCP 为集群提供了 init-actions,我现在将其用于解决方案,但我倾向于将集群用于多个 Spark 作业,现在我必须从所有作业中累积所有 python 包并在集群初始化时安装它们。所以从技术上讲,这是可行的,但似乎是 Scala/Java 通过创建具有捆绑依赖项的 fat-jars 提供的一个糟糕的替代品

标签: python python-3.x numpy apache-spark pyspark


【解决方案1】:

更新:有一个有凝聚力的 repo,其中包含一个非常出色的示例项目。你应该看看,特别是如果我下面的例子不适合你。回购在这里:https://github.com/massmutual/sample-pyspark-application 并包括这个在 YARN 上运行的示例: https://github.com/massmutual/sample-pyspark-application/blob/master/setup-and-submit.sh 期望您首先导出几个环境变量。 (我提供的值是特定于 EMR 的,因此您的值可能会有所不同。)

export HADOOP_CONF_DIR="/etc/hadoop/conf"
export PYTHON="/usr/bin/python3"
export SPARK_HOME="/usr/lib/spark"
export PATH="$SPARK_HOME/bin:$PATH"

这里提到:I can't seem to get --py-files on Spark to work 有必要使用 virtualenv 之类的东西(或者 conda 可能会工作)以避免遇到与 Python 包(例如 Numpy)的 C 库编译相关的问题,这些包依赖于底层硬件架构,无法成功移植到由于依赖项和/或任务节点中的硬链接可能与主节点实例具有不同的硬件,因此集群中的其他机器。

这里讨论了 --archives 和 --py-files 之间的一些区别:Shipping and using virtualenv in a pyspark job

我建议使用 --archives 和 virtualenv 来提供包含包依赖项的压缩文件,以避免我上面提到的一些问题。

例如,在 Amazon Elastic Map Reduce (EMR) 集群中,当 ssh 进入主实例时,我能够成功地使用 spark-submit 从 virtualenv 环境中执行测试 python 脚本,如下所示:

pip-3.4 freeze | egrep -v sagemaker > requirements.txt
# Above line is just in case you want to port installed packages across environments.
virtualenv -p python3 spark_env3
virtualenv -p python3 --relocatable spark_env3
source spark_env3/bin/activate
sudo pip-3.4 install -U pandas boto3 findspark jaydebeapi
# Note that the above libraries weren't required for the test script, but I'm showing how you can add additional dependencies if needed.
sudo pip-3.4 install -r requirements.txt
# The above line is just to show how you can load from a requirements file if needed.
cd spark_env3
# We must cd into the directory before we zip it for Spark to find the resources. 
zip -r ../spark_env3_inside.zip *
# Be sure to cd back out after building the zip file. 
cd ..

PYSPARK_PYTHON=./spark_env3/bin/python3 spark-submit \ 
  --conf spark.yarn.appMasterEnv.PYSPARK_PYTHON=./spark_env3/bin/python3 \
  --master yarn-cluster \
  --archives /home/hadoop/spark_env3_inside.zip#spark_env3 \
  test_spark.py

请注意,上面最后一行末尾附近的主题标签不是评论。它是 spark-submit 的指令,如下所述:Upload zip file using --archives option of spark-submit on yarn

我正在运行的测试脚本的来源来自这篇关于使用 conda 代替 virtualenv 来运行 pyspark 作业的文章:http://quasiben.github.io/blog/2016/4/15/conda-spark/

并包含 test_spark.py 脚本的代码:

# test_spark.py
import os
import sys
from pyspark import SparkContext
from pyspark import SparkConf

conf = SparkConf()
conf.setAppName("get-hosts")

sc = SparkContext(conf=conf)

def noop(x):
    import socket
    import sys
    return socket.gethostname() + ' '.join(sys.path) + ' '.join(os.environ)

rdd = sc.parallelize(range(1000), 100)
hosts = rdd.map(noop).distinct().collect()
print(hosts)

如果你想了解一些关于使用 virtualenv 执行 pyspark 作业的背景信息,正如 @Mariusz 已经提到的,这篇博文中有一个有用的例子:https://henning.kropponline.de/2016/09/17/running-pyspark-with-virtualenv/(尽管它没有解释我的一些微妙之处用我提供的其他链接进行了澄清)。

此处提供的答案帖子中还有一个附加示例:Elephas not loaded in PySpark: No module named elephas.spark_model

这里还有另一个示例:https://community.hortonworks.com/articles/104947/using-virtualenv-with-pyspark.html

【讨论】:

  • 谢谢。我已经有一段时间没有做这个了,但这似乎是我一直在寻找的
【解决方案2】:

如果您切换到 virtualenv,您可以轻松实现这一点。在这个环境中,您需要安装所有必要的要求,而不是压缩它并使用--archives 传递。这是一篇描述细节的好文章:https://henning.kropponline.de/2016/09/17/running-pyspark-with-virtualenv/

【讨论】:

  • 这并没有解决具有某些 C 绑定的模块的更精细问题。虽然文章引用了 numpy(其中一个特殊库),但它并未处理 ImportError OP 遇到的问题。
猜你喜欢
  • 1970-01-01
  • 2018-09-29
  • 2012-06-29
  • 2021-04-10
  • 2020-10-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-02-08
相关资源
最近更新 更多