如何在Apache Spark Dataset中从现有列生成多条新数据行
Apache Spark 多账号字段拆分为多行实现方案
你要实现的是典型的宽表转长表操作,核心是将单条记录中的3个账号字段展开为多条记录,使用Spark内置的explode行展开函数即可实现,以下是不同场景下的实现代码:
1. DataFrame API 实现(通用,支持Scala/PySpark)
Scala 版本
import org.apache.spark.sql.functions.{explode, array, col} val resultDF = inputDF .select( // 将三个账号字段打包为数组后展开为多行 explode(array("accountid", "payingaccountid", "billedaccountid")).alias("accountid"), col("startdate"), col("enddate") )
PySpark 版本
from pyspark.sql.functions import explode, array result_df = input_df.select( explode(array("accountid", "payingaccountid", "billedaccountid")).alias("accountid"), "startdate", "enddate" )
2. 强类型 Dataset 实现(Scala 专属)
如果你使用强类型Dataset开发,可以直接用flatMap算子实现:
// 定义输入、输出样例类 case class InputRow( accountid: String, payingaccountid: String, billedaccountid: String, startdate: java.sql.Timestamp, enddate: java.sql.Timestamp ) case class OutputRow( accountid: String, startdate: java.sql.Timestamp, enddate: java.sql.Timestamp ) val resultDS: Dataset[OutputRow] = inputDS.flatMap(row => Seq( OutputRow(row.accountid, row.startdate, row.enddate), OutputRow(row.payingaccountid, row.startdate, row.enddate), OutputRow(row.billedaccountid, row.startdate, row.enddate) ))
3. Spark SQL 实现
如果偏好SQL语法,可以用如下写法:
-- 先将输入数据集注册为临时视图 CREATE OR REPLACE TEMP VIEW input_view AS SELECT * FROM 你的输入表名; SELECT explode(array(accountid, payingaccountid, billedaccountid)) AS accountid, startdate, enddate FROM input_view;
注意事项
- 如果三个账号字段存在NULL值,
explode会自动跳过NULL对应的行;如果需要保留NULL值对应的行,替换explode为explode_outer即可 - 以上代码会直接沿用原行的startdate、enddate值,完全匹配需求
内容的提问来源于stack exchange,提问作者aarigoni
相关产品推荐
相关产品推荐

