【问题标题】:Spark/scala - can we create new columns from an existing column value in a dataframeSpark/scala - 我们可以从数据框中的现有列值创建新列吗
【发布时间】:2018-10-15 09:19:38
【问题描述】:

我想看看我们是否可以使用 spark/scala 从 dataFrame 中的一列中的值创建新列。 我有一个包含以下数据的数据框

df.show()

+---+-----------------------+
|id |allvals                |
+---+-----------------------+
|1  |col1,val11|col3,val31  |
|3  |col3,val33|col1,val13  |
|2  |col2,val22             |
+---+-----------------------+

在上面的数据中 col1/col2/col3 是列名后跟它的值。列名和值由, 分隔。每组由|分隔。

现在,我想实现这样的目标

+---+----------------------+------+------+------+
|id |allvals               |col1  |col2  |col3  |
+---+----------------------+------+------+------+
|1  |col1,val11|col3,val31 |val11 |null  |val31 |
|3  |col3,val33|col1,val13 |val13 |null  |val13 |
|2  |col2,val22            |null  |val22 |null  |
+---+----------------------+------+------+------+

感谢任何帮助。

【问题讨论】:

    标签: scala apache-spark spark-dataframe


    【解决方案1】:

    您可以使用udf 将列转换为Map

    import org.apache.spark.sql.functions._
    import spark.implicits._
    
    val df = Seq(
      (1, "col1,val11|col3,val31"), (2, "col3,val33|col3,val13"), (2, "col2,val22")
    ).toDF("id", "allvals")
    
    val to_map = udf((s: String) => s.split('|').collect { _.split(",") match {
      case Array(k, v) => (k, v)
    }}.toMap )
    
    val dfWithMap = df.withColumn("allvalsmap", to_map($"allvals"))
    val keys = dfWithMap.select($"allvalsmap").as[Map[String, String]].flatMap(_.keys.toSeq).distinct.collect
    
    keys.foldLeft(dfWithMap)((df, k) => df.withColumn(k, $"allvalsmap".getItem(k))).drop("allvalsmap").show
    // +---+--------------------+-----+-----+-----+
    // | id|             allvals| col3| col1| col2|
    // +---+--------------------+-----+-----+-----+
    // |  1|col1,val11|col3,v...|val31|val11| null|
    // |  2|col3,val33|col3,v...|val13| null| null|
    // |  2|          col2,val22| null| null|val22|
    // +---+--------------------+-----+-----+-----+
    

    灵感来自this answeruser6910411

    【讨论】:

    • 谢谢!这很好。我在 .as[Map[String,String]]... scala> val keys = dfWithMap.select($"allvalsmap").as[Map[String, String]].flatMap(.keys .toSeq).distinct.collect :40: 错误:无法找到存储在数据集中的类型的编码器。通过导入 spark.implicits 支持原始类型(Int、String 等)和产品类型(案例类)。 将在未来的版本中添加对序列化其他类型的支持。
    • 你必须`import spark.implicits._`,其中sparkSparkSession
    • 我导入了,但还是有问题。
    • 我无法重现。您使用哪个 Spark 版本?
    • 我使用的是 spark 2.2 和 scala 2.11
    【解决方案2】:

    您可以使用splitexplodegroupBy/pivot/agg 转换DataFrame,如下所示:

    val df = Seq(
      (1, "col1,val11|col3,val31"),
      (2, "col3,val33|col1,val13"),
      (3, "col2,val22")
    ).toDF("id", "allvals")
    
    import org.apache.spark.sql.functions._
    
    df.withColumn("temp", split($"allvals", "\\|")).
      withColumn("temp", explode($"temp")).
      withColumn("temp", split($"temp", ",")).
      select($"id", $"allvals", $"temp".getItem(0).as("k"), $"temp".getItem(1).as("v")).
      groupBy($"id", $"allvals").pivot("k").agg(first($"v"))
    
    // +---+---------------------+-----+-----+-----+
    // |id |allvals              |col1 |col2 |col3 |
    // +---+---------------------+-----+-----+-----+
    // |1  |col1,val11|col3,val31|val11|null |val31|
    // |3  |col2,val22           |null |val22|null |
    // |2  |col3,val33|col1,val13|val13|null |val33|
    // +---+---------------------+-----+-----+-----+
    

    【讨论】:

    • 谢谢利奥!更接近我的实际需求。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多