【问题标题】:How to create a pyspark udf, calling a class function from another class function in the same file?如何创建一个pyspark udf,从同一个文件中的另一个类函数调用一个类函数?
【发布时间】:2020-07-03 09:18:54
【问题描述】:

我正在基于类的视图中创建一个 pyspark udf,并且我在另一个基于类的视图中拥有我想要调用的函数,它们都在同一个文件中 (api.py),但是当我检查时结果数据框的内容,我得到这个错误:

ModuleNotFoundError: No module named 'api'

我不明白为什么会发生这种情况,我尝试在 pyspark 控制台中执行类似的代码,并且效果很好。 here 提出了类似的问题,但不同之处在于我尝试在同一个文件中执行此操作。

这是我的一段完整代码: api.py

class TextMiningMethods():
    def clean_tweet(self,tweet):
        '''
        some logic here
        '''
        return "Hello: "+tweet


class BigDataViewSet(TextMiningMethods,viewsets.ViewSet):

    @action(methods=['post'], detail=False)
    def word_cloud(self, request, *args, **kwargs): 
        '''
        some previous logic here
        '''
        spark=SparkSession \
            .builder \
            .master("spark://"+SPARK_WORKERS) \
            .appName('word_cloud') \
            .config("spark.executor.memory", '2g') \
            .config('spark.executor.cores', '2') \
            .config('spark.cores.max', '2') \
            .config("spark.driver.memory",'2g') \
            .getOrCreate()

        sc.sparkContext.addPyFile('path/to/udfFile.py')
        cols = ['text']
        rows = []

        for tweet_account_index, tweet_account_data in enumerate(tweets_list):

            tweet_data_aux_pandas_df = pd.Series(tweet_account_data['tweet']).dropna()
            for tweet_index,tweet in enumerate(tweet_data_aux_pandas_df):
                row= [tweet['text']]
                rows.append(row)

        # Create a Pandas Dataframe of tweets
        tweet_pandas_df = pd.DataFrame(rows, columns = cols)

        schema = StructType([
            StructField("text", StringType(),True)
        ])

        # Converts to Spark DataFrame
        df = spark.createDataFrame(tweet_pandas_df,schema=schema)
        clean_tweet_udf = udf(TextMiningMethods().clean_tweet, StringType())
        clean_tweet_df = df.withColumn("clean_tweet", clean_tweet_udf(df["text"]))
        clean_tweet_df.show()   # This line produces the error

pyspark 中的类似测试效果很好

from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from pyspark.sql.functions import udf
def clean_tweet(name):
    return "This is " + name

schema = StructType([StructField("Id", IntegerType(),True),StructField("tweet", StringType(),True)])

data = [[ 1, "tweet 1"],[2,"tweet 2"],[3,"tweet 3"]]
df = spark.createDataFrame(data,schema=schema)

clean_tweet_udf = udf(clean_tweet,StringType())
clean_tweet_df = df.withColumn("clean_tweet", clean_tweet_udf(df["tweet"]))
clean_tweet_df.show()

以下是我的问题:

  1. 此错误与什么有关?我该如何解决?
  2. 当您使用基于类的视图时,创建 pyspark udf 的正确方法是什么?在您将调用它们的同一文件中编写将用作 pyspark udf 的函数是错误的做法吗? (在我的情况下,我所有的 api 端点都使用 django rest 框架)

任何帮助将不胜感激,在此先感谢

更新:

这个link 和这个link 解释了如何使用SparkContext 和pyspark 一起使用自定义类,但我的情况不是SparkSession,但我使用了这个:

sc.sparkContext.addPyFile('path/to/udfFile.py')

问题是我在我为数据框创建 udf 函数的同一个文件中定义了我有函数用作 pyspark udf 的类(如我的代码中所示)。 当 addPyFile() 的路径在同一代码中时,我找不到如何实现该行为。尽管如此,我还是移动了我的代码并关注了these steps(这是我修复的另一个错误):

  • 创建一个名为udf 的新文件夹
  • 创建一个新的空__ini__.py 文件,将目录制作成一个包。
  • 并为我的 udf 函数创建一个 file.py。
core/
    udf/
    ├── __init__.py
    ├── __pycache__
    └── pyspark_udf.py
    api/
    ├── admin.py
    ├── api.py
    ├── apps.py
    ├── __init__.py

在这个文件中,我尝试在函数开头或函数内部导入依赖项。在所有情况下,我都会收到ModuleNotFoundError: No module named 'udf'

pyspark_udf.py

import re
import string
import unidecode
from nltk.corpus import stopwords

class TextMiningMethods():
    """docstring for TextMiningMethods"""
    def clean_tweet(self,tweet):
        # some logic here

我已经尝试了所有这些,在我的 api.py 文件的开头

from udf.pyspark_udf import TextMiningMethods

# or

from udf.pyspark_udf import *

在 word_cloud 函数内部

class BigDataViewSet(viewsets.ViewSet):
    def word_cloud(self, request, *args, **kwargs):
        from udf.pyspark_udf import TextMiningMethods

在 python 调试器中,这一行有效:

from udf.pyspark_udf import TextMiningMethods

但是当我显示数据框时,我收到错误:

clean_tweet_df.show()

ModuleNotFoundError: No module named 'udf'

很明显,原来的问题改成了另一个问题,现在我的问题与这个question 更相关,但是我还没有找到一个令人满意的方法来导入文件并创建一个 pyspark udf 从另一个类调用类函数功能。

我错过了什么?

【问题讨论】:

  • 什么是“基于类的视图”?
  • @Vitaliy 下面的链接解释了 django 中基于类的视图 docs.djangoproject.com/en/3.0/topics/class-based-views 是什么
  • @user10938362 我用我尝试过的所有东西更新了我的答案,从你提供给我的链接开始,这是相似但不同的情况
  • 我会尝试不同的方法:你能从 CBV 中提取逻辑并使其免费 django 吗?您的代码示例提到了文本挖掘,所以我猜核心功能 tjat 与托管无关(事实上它是由 Web 服务提供的)。我什至会做它作为鉴别诊断的实验。

标签: django python-3.x apache-spark pyspark user-defined-functions


【解决方案1】:

经过不同的尝试,我无法通过引用 addPyFile() 路径中的方法找到解决方案,该方法位于我创建 udf 的同一文件中(我想知道这是否是一种不好的做法) 或在另一个文件中,技术上 addPyFile(path) 文档说:

为将来要在此 SparkContext 上执行的所有任务添加 .py 或 .zip 依赖项。传递的路径可以是本地文件、HDFS(或其他 Hadoop 支持的文件系统)中的文件,也可以是 HTTP、HTTPS 或 FTP URI。

所以我提到的应该是可能的。基于此,我不得不使用这个 solution 并从它的最高级别压缩所有 udf 文件夹:

zip -r udf.zip udf

另外,在pyspark_udf.py 中,我必须如下导入我的依赖项以避免problem

class TextMiningMethods():
    """docstring for TextMiningMethods"""
    def clean_tweet(self,tweet):
        import re
        import string
        import unidecode
        from nltk.corpus import stopwords

代替:

import re
import string
import unidecode
from nltk.corpus import stopwords

class TextMiningMethods():
    """docstring for TextMiningMethods"""
    def clean_tweet(self,tweet):

然后,这条线终于奏效了:

clean_tweet_df.show()

我希望这对其他人有用

【讨论】:

    【解决方案2】:

    谢谢!你的方法对我有用。

    只是为了澄清我的步骤:

    • 用__init__.py 和pyspark_udfs.py 制作了一个udf 模块
    • 先制作一个 bash 文件来压缩 udfs,然后在顶层运行我的文件:

    runner.sh

    echo "zipping udfs..."
    zip -r udf.zip udf
    echo "udfs zipped"
    
    echo "running script..."
    /opt/conda/bin/python runner.py
    echo "script ended."
    
    • 在实际代码中,我从udf.pyspark_udfs 模块导入了我的udfs,并在我需要的python 函数中初始化了我的udfs,如下所示:
    
        def _produce_period_statistics(self, df: pyspark.sql.DataFrame, period: str) -> pyspark.sql.DataFrame:
    
            """ Produces basic and trend statistics based on user visits."""
    
            # udfs
            get_hist_vals_udf = F.udf(lambda array, bins, _range: get_histogram_values(array, bins, _range), ArrayType(IntegerType()))
            get_hist_edges_udf = F.udf(lambda array, bins, _range: get_histogram_edges(array, bins, _range), ArrayType(FloatType()))
            get_mean_udf = F.udf(get_mean, FloatType())
            get_std_udf = F.udf(get_std, FloatType())
            get_lr_coefs_udf = F.udf(lambda bar_height, bar_edges, hist_upper: get_linear_regression_coeffs(bar_height, bar_edges, hist_upper), StringType())
       ...
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-01-09
      • 1970-01-01
      • 2017-08-10
      • 1970-01-01
      • 2018-07-30
      • 1970-01-01
      相关资源
      最近更新 更多