pyspark returns a no module named error for a custom module

Viewed 5779

I would like to import a .py file that contains some modules. I have saved the files init.py and util_func.py under this folder:

/usr/local/lib/python3.4/site-packages/myutil

The util_func.py contains all the modules that i would like to use. I also need to create a pyspark udf so I can use it to transform my dataframe. My code looks like this:

import myutil
from myutil import util_func
myudf = pyspark.sql.functions.udf(util_func.ConvString, StringType())

somewhere down the code, I am using this to convert one of the columns in my dataframe:

df = df.withColumn("newcol", myudf(df["oldcol"]))

then I am trying to see if it converts it my using:

df.head()

It fails with an error "No module named myutil".

I am able to bring up the functions within ipython. Somehow the pyspark engined does not see the module. Any idea how to make sure that the pyspark engine picks up the module?

3 Answers

Sorry for hijack the thread. I want to reply to @rouge-one comment but I dont have enough reputation to do it

I'm having the same problem with OP but this time the module is not a single py file but the annoy spotify package in Python https://github.com/spotify/annoy/tree/master/annoy

I tried sc.addPyFile('venv.zip') and added --archives ./venv.zip#PYTHON \ in the spark-submit file but it still threw the same error message

I can still use from annoy import AnnoyIndex in the spark submit file but everytime I try to import it in the udf like this

    schema = ArrayType(StructType([
        StructField("char", IntegerType(), False),
        StructField("count", IntegerType(), False)
    ]))

    f= 128

    def return_candidate(x):
      from annoy import AnnoyIndex
      from pyspark import SparkFiles
      annoy = AnnoyIndex(f)
      annoy.load(SparkFiles.get("annoy.ann"))
      neighbor = 5
      annoy_object = annoy.get_nns_by_item(x,n = neighbor, include_distances=True)
      return annoy_object


    return_candidate_udf = udf(lambda y: return_candidate(y), schema )
inter4 =inter3.select('*',return_candidate_udf('annoy_id').alias('annoy_candidate_list'))

I found the point! Spark UDF uses another executor when you have a problem like yours, the environment variables are different!

My case, I was developing, debugging and testing on Zeppelin and it has two different interpreters for Python and Spark! When I install the libs in the terminal, I could use the functions normally but on UDF not!

Solution: Just set the same environment for driver and executor, PYSPARK_DRIVER_PYTHON and PYSPARK_PYTHON

Related