Spark DataFrame列的合并与拆分

版本说明:Spark-2.3.0

使用Spark SQL在对数据进行处理的过程中,可能会遇到对一列数据拆分为多列,或者把多列数据合并为一列。这里记录一下目前想到的对DataFrame列数据进行合并和拆分的几种方法。

1 DataFrame列数据的合并
例如:我们有如下数据,想要将三列数据合并为一列,并以“,”分割

1+----+---+-----------+ 2|name|age| phone| 3+----+---+-----------+ 4|Ming| 20|15552211521| 5|hong| 19|13287994007| 6| zhi| 21|15552211523| 7+----+---+-----------+

1.1 使用map方法重写

使用map方法重写就是将DataFrame使用map取值之后,然后使用toSeq方法转成Seq格式,最后使用Seq的foldLeft方法拼接数据,并返回,如下所示:

1//方法1:利用map重写 2 val separator = "," 3 df.map(_.toSeq.foldLeft("")(_ + separator + _).substring(1)).show() 4 5 /** 6 * +-------------------+ 7 * | value| 8 * +-------------------+ 9 * |Ming,20,15552211521| 10 * |hong,19,13287994007| 11 * | zhi,21,15552211523| 12 * +-------------------+ 13 */

1.2 使用内置函数concat_ws

合并多列数据也可以使用SparkSQL的内置函数concat_ws()

1//方法2: 使用内置函数 concat_ws 2 import org.apache.spark.sql.functions._ 3 df.select(concat_ws(separator, $"name", $"age", $"phone").cast(StringType).as("value")).show() 4 5 /** 6 * +-------------------+ 7 * | value| 8 * +-------------------+ 9 * |Ming,20,15552211521| 10 * |hong,19,13287994007| 11 * | zhi,21,15552211523| 12 * +-------------------+ 13 */

1.3 使用自定义UDF函数

自己编写UDF函数,实现多列合并

1 //方法3:使用自定义UDF函数 2 3 // 编写udf函数 4 def mergeCols(row: Row): String = { 5 row.toSeq.foldLeft("")(_ + separator + _).substring(1) 6 } 7 8 val mergeColsUDF = udf(mergeCols _) 9 df.select(mergeColsUDF(struct($"name", $"age", $"phone")).as("value")).show()

完整代码:

1import org.apache.spark.sql.{Row, SparkSession} 2import org.apache.spark.sql.types.StringType 3 4/** 5 * Created by shirukai on 2018/9/12 6 * DataFrame 合并列 7 */ 8object MergeColsTest { 9 def main(args: Array[String]): Unit = { 10 val spark = SparkSession 11 .builder() 12 .appName(this.getClass.getSimpleName) 13 .master("local") 14 .getOrCreate() 15 16 //从内存创建一组DataFrame数据 17 import spark.implicits._ 18 val df = Seq(("Ming", 20, 15552211521L), ("hong", 19, 13287994007L), ("zhi", 21, 15552211523L)) 19 .toDF("name", "age", "phone") 20 df.show() 21 /** 22 * +----+---+-----------+ 23 * |name|age| phone| 24 * +----+---+-----------+ 25 * |Ming| 20|15552211521| 26 * |hong| 19|13287994007| 27 * | zhi| 21|15552211523| 28 * +----+---+-----------+ 29 */ 30 //方法1:利用map重写 31 val separator = "," 32 df.map(_.toSeq.foldLeft("")(_ + separator + _).substring(1)).show() 33 34 /** 35 * +-------------------+ 36 * | value| 37 * +-------------------+ 38 * |Ming,20,15552211521| 39 * |hong,19,13287994007| 40 * | zhi,21,15552211523| 41 * +-------------------+ 42 */ 43 //方法2: 使用内置函数 concat_ws 44 import org.apache.spark.sql.functions._ 45 df.select(concat_ws(separator, $"name", $"age", $"phone").cast(StringType).as("value")).show() 46 47 /** 48 * +-------------------+ 49 * | value| 50 * +-------------------+ 51 * |Ming,20,15552211521| 52 * |hong,19,13287994007| 53 * | zhi,21,15552211523| 54 * +-------------------+ 55 */ 56 //方法3:使用自定义UDF函数 57 58 // 编写udf函数 59 def mergeCols(row: Row): String = { 60 row.toSeq.foldLeft("")(_ + separator + _).substring(1) 61 } 62 63 val mergeColsUDF = udf(mergeCols _) 64 df.select(mergeColsUDF(struct($"name", $"age", $"phone")).as("value")).show() 65 66 /** 67 * /** 68 * * +-------------------+ 69 * * | value| 70 * * +-------------------+ 71 * * |Ming,20,15552211521| 72 * * |hong,19,13287994007| 73 * * | zhi,21,15552211523| 74 * * +-------------------+ 75 **/ 76 */ 77 } 78}

2 DataFrame列数据的拆分

上面我们将DataFrame的多列数据合并为一列如下所示,有时候我们也需要将单列数据,以某种拆分规则,拆分为多列。下面提供几种将一列拆分为多列的方法。

1+-------------------+ 2| value| 3+-------------------+ 4|Ming,20,15552211521| 5|hong,19,13287994007| 6| zhi,21,15552211523| 7+-------------------+

2.1 使用内置函数split,然后遍历添加列

该方法,先利用内置函数split将单列的数据拆分,然后遍历使用getItem(角标)方法获取拆分后的数据,依次使用withColumn方法添加新列,代码如下所示:

1 //方法1: 使用内置函数split,然后遍历添加列 2 val separator = "," 3 lazy val first = df.first() 4 5 val numAttrs = first.toString().split(separator).length 6 val attrs = Array.tabulate(numAttrs)(n => "col_" + n) 7 //按指定分隔符拆分value列,生成splitCols列 8 var newDF = df.withColumn("splitCols", split($"value", separator)) 9 attrs.zipWithIndex.foreach(x => { 10 newDF = newDF.withColumn(x._1, $"splitCols".getItem(x._2)) 11 }) 12 newDF.show() 13 /** 14 * +-------------------+--------------------+-----+-----+-----------+ 15 * | value| splitCols|col_0|col_1| col_2| 16 * +-------------------+--------------------+-----+-----+-----------+ 17 * |Ming,20,15552211521|[Ming, 20, 155522...| Ming| 20|15552211521| 18 * |hong,19,13287994007|[hong, 19, 132879...| hong| 19|13287994007| 19 * | zhi,21,15552211523|[zhi, 21, 1555221...| zhi| 21|15552211523| 20 * +-------------------+--------------------+-----+-----+-----------+

2.2 使用UDF函数创建多列数据,然后合并
该方法是使用udf函数,生成多个列,然后合并到原来的数据。该方法参考了VectorDisassembler(与spark ml官网提供的VectorAssembler相反),这是一个第三方的spark ml向量拆分算法,该方法github地址:https://github.com/jamesbconner/VectorDisassembler。代码如下所示:

1//方法2:使用udf函数创建多列,然后合并 2 val attributes: Array[Attribute] = { 3 val numAttrs = first.toString().split(separator).length 4 //生成attributes 5 Array.tabulate(numAttrs)(i => NumericAttribute.defaultAttr.withName("value" + "_" + i)) 6 } 7 //创建多列数据 8 val fieldCols = attributes.zipWithIndex.map(x => { 9 val assembleFunc = udf { 10 str: String => 11 str.split(separator)(x._2) 12 } 13 assembleFunc(df("value").cast(StringType)).as(x._1.name.get, x._1.toMetadata()) 14 }) 15 //合并数据 16 df.select(col("*") +: fieldCols: _*).show() 17 18 /** 19 * +-------------------+-------+-------+-----------+ 20 * | value|value_0|value_1| value_2| 21 * +-------------------+-------+-------+-----------+ 22 * |Ming,20,15552211521| Ming| 20|15552211521| 23 * |hong,19,13287994007| hong| 19|13287994007| 24 * | zhi,21,15552211523| zhi| 21|15552211523| 25 * +-------------------+-------+-------+-----------+ 26 */

完整代码:

1import org.apache.spark.ml.attribute.{Attribute, NumericAttribute} 2import org.apache.spark.sql.SparkSession 3import org.apache.spark.sql.types.StringType 4 5/** 6 * Created by shirukai on 2018/9/12 7 * 拆分列 8 */ 9object SplitColTest { 10 def main(args: Array[String]): Unit = { 11 val spark = SparkSession 12 .builder() 13 .appName(this.getClass.getSimpleName) 14 .master("local") 15 .getOrCreate() 16 17 //从内存中创建DataFrame 18 import spark.implicits._ 19 val df = Seq("Ming,20,15552211521", "hong,19,13287994007", "zhi,21,15552211523") 20 .toDF("value") 21 df.show() 22 23 /** 24 * +-------------------+ 25 * | value| 26 * +-------------------+ 27 * |Ming,20,15552211521| 28 * |hong,19,13287994007| 29 * | zhi,21,15552211523| 30 * +-------------------+ 31 */ 32 33 import org.apache.spark.sql.functions._ 34 //方法1: 使用内置函数split,然后遍历添加列 35 val separator = "," 36 lazy val first = df.first() 37 38 val numAttrs = first.toString().split(separator).length 39 val attrs = Array.tabulate(numAttrs)(n => "col_" + n) 40 //按指定分隔符拆分value列,生成splitCols列 41 var newDF = df.withColumn("splitCols", split($"value", separator)) 42 attrs.zipWithIndex.foreach(x => { 43 newDF = newDF.withColumn(x._1, $"splitCols".getItem(x._2)) 44 }) 45 newDF.show() 46 47 /** 48 * +-------------------+--------------------+-----+-----+-----------+ 49 * | value| splitCols|col_0|col_1| col_2| 50 * +-------------------+--------------------+-----+-----+-----------+ 51 * |Ming,20,15552211521|[Ming, 20, 155522...| Ming| 20|15552211521| 52 * |hong,19,13287994007|[hong, 19, 132879...| hong| 19|13287994007| 53 * | zhi,21,15552211523|[zhi, 21, 1555221...| zhi| 21|15552211523| 54 * +-------------------+--------------------+-----+-----+-----------+ 55 */ 56 57 //方法2:使用udf函数创建多列,然后合并 58 val attributes: Array[Attribute] = { 59 val numAttrs = first.toString().split(separator).length 60 //生成attributes 61 Array.tabulate(numAttrs)(i => NumericAttribute.defaultAttr.withName("value" + "_" + i)) 62 } 63 //创建多列数据 64 val fieldCols = attributes.zipWithIndex.map(x => { 65 val assembleFunc = udf { 66 str: String => 67 str.split(separator)(x._2) 68 } 69 assembleFunc(df("value").cast(StringType)).as(x._1.name.get, x._1.toMetadata()) 70 }) 71 //合并数据 72 df.select(col("*") +: fieldCols: _*).show() 73 74 /** 75 * +-------------------+-------+-------+-----------+ 76 * | value|value_0|value_1| value_2| 77 * +-------------------+-------+-------+-----------+ 78 * |Ming,20,15552211521| Ming| 20|15552211521| 79 * |hong,19,13287994007| hong| 19|13287994007| 80 * | zhi,21,15552211523| zhi| 21|15552211523| 81 * +-------------------+-------+-------+-----------+ 82 */ 83 } 84}
点赞
收藏

评论区

加载中...

相关推荐

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

MySQL查询结果集字符串操作之多行合并与单行分割

前言我们在做项目写sql语句的时候,是否会遇到这样的场景,就是需要把查询出来的多列,按照字符串分割合并成一列显示,或者把存在数据库里面用逗号分隔的一列,查询分成多列呢,常见场景有,文章标签,需要吧查询多个标签合并成一列,等,需要怎么去实现呢,这就涉及到MySQL的字符串操作groupconcat场景再现我想把查询多列数据合并成一列显示用逗号分隔

Excel数据转化为sql脚本

在实际项目开发中,有时会遇到客户让我们把大量Excel数据导入数据库的情况。这时我们就可以通过将Excel数据转化为sql脚本来批量导入数据库。1在数据前插入一列单元格,用来拼写sql语句。 具体写法:"insertintot\_student(id,name,age,class)value("&B2&",'"&C2&"',"&D2&"

95%的人都不知道 MySQL还有索引管理与执行计划

1.1索引的介绍  索引是对数据库表中一列或多列的值进行排序的一种结构,使用索引可快速访问数据库表中的特定信息。如果想按特定职员的姓来查找他或她,则与在表中搜索所有的行相比,索引有助于更快地获取信息。  索引的一个主要目的就是加快检索表中数据的方法,亦即能协助信息搜索者尽快的找到符合限制条件的记录ID的辅助数据结构。!fi

Python之DataFrame更改列名及重排列顺序

日常在处理数据的时候,经常需要对dataframe进行重排,只取其中几列或者更改列名等操作;有两个相似的方法reindex和rename,与此记录一下常见的用法,并标注一下区别:rename:重命名,就是对col列进行命名的修改,他只改变col的名字,相当于起了个别名,原来叫col1,以后叫col2,inplaceTrue,用来保存更改,即更改了原

mysql——GROUP BY和HAVING

GROUPBY语法可以根据给定数据列的每个成员对查询结果进行分组统计,最终得到一个分组汇总表。select子句中的列名必须为分组列或列函数,列函数对于groupby子句定义的每个组返回一个结果。某个员工信息表结构和数据如下:  id  name  dept  salary  edlevel     hiredate   1  张