【问题标题】:PySpark - Pass list as parameter to UDFPySpark - 将列表作为参数传递给 UDF
【发布时间】:2018-06-03 01:16:03
【问题描述】:

我需要将一个列表传递给一个 UDF,该列表将确定距离的分数/类别。目前,我将所有距离硬编码为第 4 分。

a= spark.createDataFrame([("A", 20), ("B", 30), ("D", 80)],["Letter", "distances"])

from pyspark.sql.functions import udf
def cate(label, feature_list):
    if feature_list == 0:
        return label[4]
label_list = ["Great", "Good", "OK", "Please Move", "Dead"]
udf_score=udf(cate, StringType())
a.withColumn("category", udf_score(label_list,a["distances"])).show(10)

当我尝试这样的事情时,我得到了这个错误。

Py4JError: An error occurred while calling z:org.apache.spark.sql.functions.col. Trace:
py4j.Py4JException: Method col([class java.util.ArrayList]) does not exist
    at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:318)
    at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:339)
    at py4j.Gateway.invoke(Gateway.java:274)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.GatewayConnection.run(GatewayConnection.java:214)
    at java.lang.Thread.run(Thread.java:745)

【问题讨论】:

    标签: python pyspark spark-dataframe user-defined-functions


    【解决方案1】:

    尝试对函数进行柯里化,以便 DataFrame 调用中的唯一参数是您希望函数对其执行操作的列的名称:

    udf_score=udf(lambda x: cate(label_list,x), StringType())
    a.withColumn("category", udf_score("distances")).show(10)
    

    【讨论】:

      【解决方案2】:
      from pyspark.sql.functions import udf, col
      
      #sample data
      a= sqlContext.createDataFrame([("A", 20), ("B", 30), ("D", 80)],["Letter", "distances"])
      label_list = ["Great", "Good", "OK", "Please Move", "Dead"]
      
      def cate(label, feature_list):
          if feature_list == 0:
              return label[4]
          else:  #you may need to add 'else' condition as well otherwise 'null' will be added in this case
              return 'I am not sure!'
      
      def udf_score(label_list):
          return udf(lambda l: cate(l, label_list))
      a.withColumn("category", udf_score(label_list)(col("distances"))).show()
      

      输出是:

      +------+---------+--------------+
      |Letter|distances|      category|
      +------+---------+--------------+
      |     A|       20|I am not sure!|
      |     B|       30|I am not sure!|
      |     D|       80|I am not sure!|
      +------+---------+--------------+
      

      【讨论】:

      • 人们说我们可以使用 pyspark.sql.functions.array() 直接将列表传递给 UDF(从 Spark 2.20 开始)。如何使用 array() 重写上面的示例。
      • 您的解决方案有问题。它总是只剩下一部分。
      【解决方案3】:

      我认为通过将列表作为变量的默认值传递可能会有所帮助

      from pyspark.sql.functions import udf, col
      
      #sample data
      a= sqlContext.createDataFrame([("A", 20), ("B", 30), ("D", 80),("E",0)],["Letter", "distances"])
      label_list = ["Great", "Good", "OK", "Please Move", "Dead"]
      
      #Passing List as Default value to a variable
      def cate( feature_list,label=label_list):
          if feature_list == 0:
              return label[4]
          else:  #you may need to add 'else' condition as well otherwise 'null' will be added in this case
              return 'I am not sure!'
      
      udfcate = udf(cate, StringType())
      
      a.withColumn("category", udfcate("distances")).show()
      

      输出:

      +------+---------+--------------+
      |Letter|distances|      category|
      +------+---------+--------------+
      |     A|       20|I am not sure!|
      |     B|       30|I am not sure!|
      |     D|       80|I am not sure!|
      |     E|        0|          Dead|
      +------+---------+--------------+
      

      【讨论】:

      • 抱歉投反对票;我觉得问题更多是关于如何确定将两个参数发送给函数,而不是默认情况下始终使用一个参数。如果 label_list 未在全局范围内定义,并且您需要动态发送该列表,则此解决方案无法应对。 ags29 和@Prem 准确地回答了这个问题。甚至我也在寻找类似的解决方案。
      • Amith.. 该列表在示例中是静态的,因此对于静态列表来说这是一个公平可行的解决方案。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-03-01
      • 2017-05-23
      • 2019-04-02
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多