【问题标题】:Passing Array to Spark Lit function将数组传递给 Spark Lit 函数
【发布时间】:2018-09-15 22:58:03
【问题描述】:

假设我有一个包含数字 1-10 的 numpy 数组 a
[1 2 3 4 5 6 7 8 9 10]

我还有一个 Spark 数据框,我想向其中添加我的 numpy 数组 a。我认为一列文字可以完成这项工作。这不起作用:

df = df.withColumn("NewColumn", F.lit(a))

不支持的文字类型类 java.util.ArrayList

但这有效:

df = df.withColumn("NewColumn", F.lit(a[0]))

怎么做?

之前的示例 DF:

col1
a b c d e f g h i j

预期结果:

col1 NewColumn
a b c d e f g h i j 1 2 3 4 5 6 7 8 9 10

【问题讨论】:

    标签: python apache-spark pyspark apache-spark-sql literals


    【解决方案1】:

    Spark 的array 中的列表理解

    a = [1,2,3,4,5,6,7,8,9,10]
    df = spark.createDataFrame([['a b c d e f g h i j '],], ['col1'])
    df = df.withColumn("NewColumn", F.array([F.lit(x) for x in a]))
    
    df.show(truncate=False)
    df.printSchema()
    #  +--------------------+-------------------------------+
    #  |col1                |NewColumn                      |
    #  +--------------------+-------------------------------+
    #  |a b c d e f g h i j |[1, 2, 3, 4, 5, 6, 7, 8, 9, 10]|
    #  +--------------------+-------------------------------+
    #  root
    #   |-- col1: string (nullable = true)
    #   |-- NewColumn: array (nullable = false)
    #   |    |-- element: integer (containsNull = false)
    

    @pault 评论 (Python 2.7)

    您可以使用map 隐藏循环:
    df.withColumn("NewColumn", F.array(map(F.lit, a)))

    @abegehr 添加了Python 3版本:

    df.withColumn("NewColumn", F.array(*map(F.lit, a)))

    Spark 的udf

    # Defining UDF
    def arrayUdf():
        return a
    callArrayUdf = F.udf(arrayUdf, T.ArrayType(T.IntegerType()))
    
    # Calling UDF
    df = df.withColumn("NewColumn", callArrayUdf())
    

    输出是一样的。

    【讨论】:

    • 我试过了,效果很好。谢谢你的回答,我现在会保持这种方式。然而,实际上,我的“a”数组有数以万计的条目,并且由于 for 循环,它的效率不是很高。有没有办法不用循环?
    • @A.R.我已经使用不需要 for 循环的 udf 函数更新了我的答案。如果答案有帮助,您可以接受并点赞
    • 你可以使用map隐藏循环:df.withColumn("NewColumn", F.array(map(F.lit, a)))
    • @pault map 不是 rdd 函数吗?此外,map 的输出既不是字符串也不是列,因此 withColumn 会抛出错误。
    • @pault,我认为这应该是 F.array(*map(F.lit, a)) 与(星)扩展运算符,因为 F.array 无法处理地图对象。
    【解决方案2】:

    在 scala API 中,我们可以使用“typedLit”函数在列中添加 Array 或 map 值。

    // 参考:https://spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.sql.functions$

    这里是添加 Array 或 Map 作为列值的示例代码。

    import org.apache.spark.sql.functions.typedLit
    
    val df1 = Seq((1, 0), (2, 3)).toDF("a", "b")
    
    df1.withColumn("seq", typedLit(Seq(1,2,3)))
        .withColumn("map", typedLit(Map(1 -> 2)))
        .show(truncate=false)
    

    // 输出

    +---+---+---------+--------+
    |a  |b  |seq      |map     |
    +---+---+---------+--------+
    |1  |0  |[1, 2, 3]|[1 -> 2]|
    |2  |3  |[1, 2, 3]|[1 -> 2]|
    +---+---+---------+--------+
    

    我希望这会有所帮助。

    【讨论】:

    • 这没有回答问题,OP 要求提供 pyspark 解决方案。
    猜你喜欢
    • 2015-09-11
    • 2022-07-28
    • 2020-07-28
    • 1970-01-01
    • 1970-01-01
    • 2021-07-01
    相关资源
    最近更新 更多