【问题标题】:Spark: map columns of a dataframe to their ID of the distinct elementsSpark:将数据框的列映射到不同元素的 ID
【发布时间】:2021-04-30 04:45:47
【问题描述】:

我有以下两列字符串类型 A 和 B 的数据框:

val df = (
    spark
    .createDataFrame(
        Seq(
            ("a1", "b1"),
            ("a1", "b2"),
            ("a1", "b2"),
            ("a2", "b3")
        )
    )
).toDF("A", "B")

我在每列的不同元素和一组整数之间创建映射

val mapColA = (
    df
    .select("A")
    .distinct
    .rdd
    .zipWithIndex
    .collectAsMap
)

val mapColB = (
    df
    .select("B")
    .distinct
    .rdd
    .zipWithIndex
    .collectAsMap
)

现在我想在数据框中创建一个新列,将这些映射应用到它们对应的列。仅对于一张地图,这将是

df.select("A").map(x=>mapColA.get(x)).show()

但是我不明白如何将每个映射应用到其对应的列并创建两个新列(例如 withColumn)。预期的结果是

val result = (
    spark
    .createDataFrame(
        Seq(
            ("a1", "b1", 1, 1),
            ("a1", "b2", 1, 2),
            ("a1", "b2", 1, 2),
            ("a2", "b3", 2, 3)
        )
    )
).toDF("A", "B", "idA", "idB")

你能帮帮我吗?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql rdd


    【解决方案1】:

    如果我理解正确的话,这可以使用dense_rank来实现:

    import org.apache.spark.sql.expressions.Window
    
    val df2 = df.withColumn("idA", dense_rank().over(Window.orderBy("A")))
                .withColumn("idB", dense_rank().over(Window.orderBy("B")))
    
    df2.show
    +---+---+---+---+
    |  A|  B|idA|idB|
    +---+---+---+---+
    | a1| b1|  1|  1|
    | a1| b2|  1|  2|
    | a1| b2|  1|  2|
    | a2| b3|  2|  3|
    +---+---+---+---+
    

    如果你想坚持原来的代码,你可以做一些修改:

    val mapColA = df.select("A").distinct().rdd.map(r=>r.getAs[String](0)).zipWithIndex.collectAsMap
    
    val mapColB = df.select("B").distinct().rdd.map(r=>r.getAs[String](0)).zipWithIndex.collectAsMap
    
    val df2 = df.map(r => (r.getAs[String](0), r.getAs[String](1), mapColA.get(r.getAs[String](0)), mapColB.get(r.getAs[String](1)))).toDF("A","B", "idA", "idB")
    
    df2.show
    +---+---+---+---+
    |  A|  B|idA|idB|
    +---+---+---+---+
    | a1| b1|  1|  2|
    | a1| b2|  1|  0|
    | a1| b2|  1|  0|
    | a2| b3|  0|  1|
    +---+---+---+---+
    

    【讨论】:

    • 非常好的建议!但是,我在更大的数据集中收到了 WARN WindowExec: No Partition Defined for Window operation! Moving all data to a single partition, this can cause serious performance degradation. 错误。我想知道是否有更好的方法来完成同样的任务
    • @Galuoises 我添加了一种基于您的原始代码的替代方法。看看有没有帮助!
    • 感谢您的代码!这正是我想要构建的
    猜你喜欢
    • 2018-09-20
    • 2017-02-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-14
    • 1970-01-01
    相关资源
    最近更新 更多