Spark Scala实现DataFrame/Parquet文件一行转多行
解决Spark DataFrame单行转多行的问题
嘿,我明白你要做的是把每条输入记录拆成两条输出行——一条带部门/年份信息(还有个和Id对应的序号),另一条带语言信息。下面用Spark Scala给你写一套完整的实现方案:
1. 先定义输入数据的Schema
首先得明确输入数据的结构,我们先定义对应的Schema,不管你是从Parquet加载还是用样例数据,这个Schema都能用:
import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ val inputSchema = StructType(Seq( StructField("Id", IntegerType, nullable = false), StructField("PersonName", StringType, nullable = false), StructField("Dept", IntegerType, nullable = false), StructField("year", IntegerType, nullable = false), StructField("Language", StringType, nullable = false) ))
2. 加载或构造输入DataFrame
如果你的数据存在Parquet文件里,直接用spark.read.parquet("path/to/your/file")加载就行。这里我用样例数据模拟你的输入:
val inputData = Seq( (1, "David", 501, 2018, "English"), (2, "Nancy", 501, 2018, "English"), (3, "Shyam", 502, 2018, "Hindi") ).toDF(inputSchema.fieldNames: _*)
3. 实现行转多行的核心逻辑
这里有两种方式,都能达到你的需求:
方式一:用flatMap直接拆分记录
这种方式最直观,把每条输入记录映射成两条输出记录:
val transformedDF = inputData.flatMap { row => val id = row.getAs[Int]("Id") val name = row.getAs[String]("PersonName") val dept = row.getAs[Int]("Dept") val year = row.getAs[Int]("year") val lang = row.getAs[String]("Language") // 生成两条记录:第一条是部门年份行,第二条是语言行 Seq( (id, name, id, dept, year), // 第三个字段是你输出里的序号,和Id一致 (id, name, lang) ) }
方式二:拆分后再合并(更易维护)
如果后续需要调整两种记录的逻辑,拆成两个DataFrame再合并会更清晰:
// 提取部门年份相关的记录 val deptYearRows = inputData.select($"Id", $"PersonName", $"Id".alias("SeqNo"), $"Dept", $"year") // 提取语言相关的记录 val languageRows = inputData.select($"Id", $"PersonName", $"Language") // 合并两个DataFrame,允许列数不同 val finalDF = deptYearRows.unionByName(languageRows, allowMissingColumns = true)
4. 查看输出结果
执行finalDF.show(false)就能看到和你期望一致的结构:
+---+----------+-----+----+----+-------+ |Id |PersonName|SeqNo|Dept|year|Language| +---+----------+-----+----+----+-------+ |1 |David |1 |501 |2018|null | |2 |Nancy |2 |501 |2018|null | |3 |Shyam |3 |502 |2018|null | |1 |David |null |null|null|English | |2 |Nancy |null |null|null|English | |3 |Shyam |null |null|null|Hindi | +---+----------+-----+----+----+-------+
如果需要输出成无null的纯文本格式,你可以用foreach或者写入文件时处理,比如:
finalDF.foreach { row => // 过滤掉null值,只打印存在的字段 val nonNullValues = row.toSeq.filter(_ != null).mkString(" ") println(nonNullValues) }
这样打印出来的结果就完全和你给出的例子一致啦!
内容的提问来源于stack exchange,提问作者Arvy
相关产品推荐
相关产品推荐

