代码之家  ›  专栏  ›  技术社区  ›  Go Erlangen

PySpark UDF返回大小可变的元组

  •  4
  • Go Erlangen  · 技术社区  · 8 年前

    我获取一个现有的数据帧,并创建一个包含元组的字段的新数据帧。自定义项用于生成此字段。例如,在这里,我获取一个源元组并修改其元素以生成一个新元组:

    udf( lambda x: tuple([2*e for e in x], ...)
    

    所面临的挑战是元组的长度事先未知,并且可以在行与行之间更改。

    根据我对阅读相关讨论的理解,要返回元组,UDF的返回类型必须声明为StructType。然而,由于返回的元组中的元素数量未知,我不能只编写如下内容:

    StructType([
        StructField("w1", IntegerType(), False),
        StructField("w2", IntegerType(), False),
        StructField("w3", IntegerType(), False)])
    

    似乎可以返回列表,但列表对我不起作用,因为我需要在输出数据帧中使用哈希对象。

    我的选择是什么?

    提前感谢

    2 回复  |  直到 7 年前
        1
  •  2
  •   zero323 little_kid_pea    7 年前

    StructType / Row 表示固定大小 product type 对象,不能用于表示可变大小的对象。

    要表示同构集合,请使用 list 作为外部类型,以及 ArrayType 作为SQL类型:

    udf(lambda x: [2*e for e in x], ArrayType(IntegerType()))
    

    或(Spark 2.2或更高版本):

    udf(lambda x: [2*e for e in x], "array<integer>")
    

    在Spark 2.4或更高版本中,您可以使用 transform

    from pyspark.sql.functions import expr
    
    expr("tranform(input_column, x -> 2 * x)")
    
        2
  •  0
  •   Aus_10    7 年前

    每次一行的新语法per-Databricks(Spark)(语法更符合Pandas UDF,这似乎是python中UDF的发展方向 https://databricks.com/blog/2017/10/30/introducing-vectorized-udfs-for-pyspark.html ):

    一次一行:

    @udf(ArrayType(IntegerType()))
    def new_tuple(x):
        return [2*e for e in x]