Spark-Scala实现:基于另一列值拆分生成多行
用Spark-Scala实现按size列拆分多行数据的需求
输入数据
| id | size |
|---|---|
| 1 | 4 |
| 2 | 2 |
需求说明
根据size列的值拆分生成多行数据:
- 若
size值为N,则拆分为N行,新的size字段值依次为1到N - 原始行的
id保持不变,对应拆分后的所有行
期望输出
| id | size |
|---|---|
| 1 | 1 |
| 1 | 2 |
| 1 | 3 |
| 1 | 4 |
| 2 | 1 |
| 2 | 2 |
Spark-Scala实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object SplitRowsBySize { def main(args: Array[String]): Unit = { // 初始化SparkSession val spark = SparkSession.builder() .appName("SplitRowsBySize") .master("local[*]") // 本地运行模式,生产环境可移除 .getOrCreate() import spark.implicits._ // 创建输入DataFrame val inputDF = Seq( (1, 4), (2, 2) ).toDF("id", "size") // 核心逻辑:用sequence生成序列,再用explode拆分行 val resultDF = inputDF.withColumn( "size", explode(sequence(lit(1), col("size"))) ) // 展示结果 resultDF.show() // 停止SparkSession spark.stop() } }
代码说明
- sequence函数:生成从1到原始
size值的整数序列,比如原始值为4时生成[1,2,3,4] - explode函数:将数组类型的列拆分成多行,每行对应数组中的一个元素,同时保留原始行的
id字段 - 全程使用Spark内置函数,无需自定义UDF,实现高效简洁
内容的提问来源于stack exchange,提问作者sunny singh
相关产品推荐
相关产品推荐

