代码之家  ›  专栏  ›  技术社区  ›  Ivan Bilan

PySpark中未加载Elephas:没有名为Elephas的模块。spark\u型号

  •  1
  • Ivan Bilan  · 技术社区  · 8 年前

    我正在尝试在集群上分发Keras培训,并使用Elephas来实现这一目标。但是,当运行Elephas文档中的基本示例时( https://github.com/maxpumperla/elephas ):

    from elephas.utils.rdd_utils import to_simple_rdd
    rdd = to_simple_rdd(sc, x_train, y_train)
    from elephas.spark_model import SparkModel
    from elephas import optimizers as elephas_optimizers
    sgd = elephas_optimizers.SGD()
    spark_model = SparkModel(sc, model, optimizer=sgd, frequency='epoch', mode='asynchronous', num_workers=2)
    spark_model.train(rdd, nb_epoch=epochs, batch_size=batch_size, verbose=1, validation_split=0.1)
    

    我得到以下错误:

     ImportError: No module named elephas.spark_model
    
    
    
    ```Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.collectAndServe.
    : org.apache.spark.SparkException: Job aborted due to stage failure: Task 1 in stage 5.0 failed 4 times, most recent failure: Lost task 1.3 in stage 5.0 (TID 58, xxxx, executor 8): org.apache.spark.api.python.PythonException: Traceback (most recent call last):
      File "/xx/xx/hadoop/yarn/local/usercache/xx/appcache/application_151xxx857247_19188/container_1512xxx247_19188_01_000009/pyspark.zip/pyspark/worker.py", line 163, in main
        func, profiler, deserializer, serializer = read_command(pickleSer, infile)
      File "/xx/xx/hadoop/yarn/local/usercache/xx/appcache/application_151xxx857247_19188/container_1512xxx247_19188_01_000009/pyspark.zip/pyspark/worker.py", line 54, in read_command
        command = serializer._read_with_length(file)
      File /yarn/local/usercache/xx/appcache/application_151xxx857247_19188/container_1512xxx247_19188_01_000009/pyspark.zip/pyspark/serializers.py", line 169, in _read_with_length
        return self.loads(obj)
      File "/yarn//local/usercache/xx/appcache/application_151xxx857247_19188/container_1512xxx247_19188_01_000009/pyspark.zip/pyspark/serializers.py", line 454, in loads
        return pickle.loads(obj)
    ImportError: No module named elephas.spark_model
    
        at org.apache.spark.api.python.PythonRunner$$anon$1.read(PythonRDD.scala:193)
        at org.apache.spark.api.python.PythonRunner$$anon$1.<init>(PythonRDD.scala:234)
        at org.apache.spark.api.python.PythonRunner.compute(PythonRDD.scala:152)
        at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:63)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:323)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:287)
        at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
        at org.apache.spark.scheduler.Task.run(Task.scala:99)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:322)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)```
    

    而且,模型实际上是创建的,我可以 print(spark_model) 我会得到这个 <elephas.spark_model.SparkModel object at 0x7efce0abfcd0> 。错误发生在 spark_model.train

    我已使用安装elephas pip2 install git+https://github.com/maxpumperla/elephas ,也许这是相关的。

    我使用PySpark 2.1.1、Keras 2.1.4和Python 2.7。 我试过用 spark提交 :

    PYSPARK_DRIVER_PYTHON=`which python` spark-submit --driver-memory 1G  filname.py
    

    也可以直接放在Jupyter笔记本中。两者都会导致相同的问题。

    有人能给我一些建议吗?这与elephas有关还是Pypark问题?

    编辑:我还上传了虚拟环境的zip文件,并在脚本中调用它:

    virtualenv spark_venv --relocatable
    cd spark_venv 
    zip -qr ../spark_venv.zip *
    
    PYSPARK_DRIVER_PYTHON=`which python` spark-submit --driver-memory 1G --py-files spark_venv.zip filename.py
    

    然后在文件中,我执行以下操作:

    sc.addPyFile("spark_venv.zip")
    

    在这个keras导入后没有任何问题,但我仍然得到 elephas 上面的错误。

    2 回复  |  直到 8 年前
        1
  •  2
  •   Ivan Bilan    8 年前

    我找到了一个解决方案,可以正确地将虚拟环境加载到主工作环境和所有从工作环境:

    virtualenv venv --relocatable
    cd venv 
    zip -qr ../venv.zip *
    
    PYSPARK_PYTHON=./SP/bin/python spark-submit --master yarn --deploy-mode cluster --conf spark.yarn.appMasterEnv.PYSPARK_PYTHON=./SP/bin/python --driver-memory 4G --archives venv.zip#SP filename.py
    

    GitHub问题中的更多详细信息: https://github.com/maxpumperla/elephas/issues/80#issuecomment-371073492

        2
  •  1
  •   addmeaning    8 年前

    您应该添加 elephas 库作为您的 spark-submit 命令

    引用官方指南:

    对于Python,可以使用 --py-files 的参数 spark提交 添加。py。拉链或。要随应用程序分发的egg文件。如果您依赖多个Python文件,我们建议将它们打包到。拉链或。鸡蛋

    Official guide