【问题标题】:How can I write dynamic explode function(to explode multiple columns) in Scala如何在 Scala 中编写动态分解函数(分解多列)
【发布时间】:2020-07-21 14:42:58
【问题描述】:

我需要编写一个动态的 Scala 类。它将接受三个参数作为输入。 input_dataframe,要分解的列列表和分隔符。考虑我有以下数据框。

DataBase     TableName       Value
dbdev        table1_name     Value1#Value2#Value3

爆炸后我期待结果如下

DataBase     TableName       Value                   Value_Exploded
dbdev        table1_name     Value1#Value2#Value3    Value1
dbdev        table1_name     Value1#Value2#Value3    Value2
dbdev        table1_name     Value1#Value2#Value3    Value3

所以我的问题是如何编写一个Scala类来实现上面的。 约束是,它必须是通用的。它可能会得到不同的数据框。并且需要分解的列(多个)需要传递。

当我只需要分解一列时,我能够做到这一点。请在下面找到-

val explodeColumnName = "Value" //column which i need to explode
val explodeColumnBy = "#" //delimiter

val explodeDF = df.select(df.col("*"), explode(split(col(explodeColumnName), s"$explodeColumnBy")).as (explodeColumnName+"_Exploded"))

但是当我需要动态分解多个列时,我失败了。前任。假设,我需要分解 Dataframe df 的 4 列。

任何帮助/建议/建议都会非常棒。

谢谢!

【问题讨论】:

  • 在选择中分解多个列将不起作用..您需要使用 withColumn 函数或尝试以下解决方案..如果它不起作用,请告诉我..

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


【解决方案1】:

检查下面的代码。

scala> val df = Seq(
     (
         "dbdev",
         "table1_name",
         "Value1#Value2#Value3",
         "Sample1#Sample2#Sample3"
    )
)
.toDF("database","tablename","value","sample")
scala> df.show(false)
+--------+-----------+--------------------+-----------------------+
|database|tablename  |value               |sample                 |
+--------+-----------+--------------------+-----------------------+
|dbdev   |table1_name|Value1#Value2#Value3|Sample1#Sample2#Sample3|
+--------+-----------+--------------------+-----------------------+

导入所需的库

scala> import org.apache.spark.sql.{Column,DataFrame}
import org.apache.spark.sql.{Column, DataFrame}

定义DFHelper 类。

注意 - 不要在DFHelper 类中使用explode 作为函数名,explode 已经在内置函数中可用,所以我使用explodeM 作为函数。

scala> implicit class DFHelper(inDF: DataFrame) {
           import inDF.sparkSession.implicits._          
            def explodeM(delimiter:String,columns:Column*): DataFrame = {
               columns.foldLeft(inDF)((indf,column) => indf
               .withColumn(column.toString,split(column,delimiter))
               .withColumn(column.toString,explode(column))
               )
           }
      }

scala> df.explodeM("#",$"value").show(false) // one column exploding
+--------+-----------+------+-----------------------+
|database|tablename  |value |sample                 |
+--------+-----------+------+-----------------------+
|dbdev   |table1_name|Value1|Sample1#Sample2#Sample3|
|dbdev   |table1_name|Value2|Sample1#Sample2#Sample3|
|dbdev   |table1_name|Value3|Sample1#Sample2#Sample3|
+--------+-----------+------+-----------------------+
scala> df.explodeM("#",$"value",$"sample").show(false) // two columns exploding
+--------+-----------+------+-------+
|database|tablename  |value |sample |
+--------+-----------+------+-------+
|dbdev   |table1_name|Value1|Sample1|
|dbdev   |table1_name|Value1|Sample2|
|dbdev   |table1_name|Value1|Sample3|
|dbdev   |table1_name|Value2|Sample1|
|dbdev   |table1_name|Value2|Sample2|
|dbdev   |table1_name|Value2|Sample3|
|dbdev   |table1_name|Value3|Sample1|
|dbdev   |table1_name|Value3|Sample2|
|dbdev   |table1_name|Value3|Sample3|
+--------+-----------+------+-------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-06-14
    • 2017-02-12
    • 2020-05-09
    • 2021-05-22
    • 1970-01-01
    • 1970-01-01
    • 2015-09-17
    相关资源
    最近更新 更多