Spark Scala:如何拆分名为rankedDF的DataFrame?
拆分Spark Scala DataFrame的Items列为多行
搞定这个需求很简单,我们可以用Spark内置的split和explode函数组合实现,直接看具体步骤:
1. 导入必要的函数
首先要导入Spark SQL的相关函数,不然没法用split和explode:
import org.apache.spark.sql.functions.{split, explode, col, when, array}
2. 核心拆分逻辑
下面的代码会把rankedDF里的Items列(逗号分隔的字符串)拆分成单独的行,同时保留其他列的所有值:
val splitItemsDF = rankedDF // 第一步:把Items列按「逗号+任意空格」分割成数组,避免拆分后Item带前导空格 .withColumn("ItemArray", split(col("Items"), ",\\s*")) // 第二步:把数组中的每个元素拆成单独的行 .withColumn("Item", explode(col("ItemArray"))) // 第三步:清理临时列和原Items列 .drop("ItemArray", "Items")
代码细节解释
split(col("Items"), ",\\s*"):用正则表达式,\\s*匹配逗号和后面的任意空格,这样原数据里类似"Womens Socks, Men Pants"的内容,拆分后会得到["Womens Socks", "Men Pants"],而不是带空格的["Womens Socks", " Men Pants"]explode(col("ItemArray")):这一步是关键,它会把数组中的每个元素生成一条新行,其他列(比如TimePeriod、TXN_HEADER_ID)的值和原行完全一致- 最后
drop掉临时生成的ItemArray和原来的Items列,得到干净的拆分结果
3. 处理特殊情况(可选)
如果你的Items列可能存在null或者空字符串的情况,可以加个判断避免生成空行:
val splitItemsDF = rankedDF .withColumn("ItemArray", // 如果Items不为空才拆分,否则返回空数组(也可以根据需求返回null) when(col("Items").isNotNull && col("Items") =!= "", split(col("Items"), ",\\s*")) .otherwise(array()) ) .withColumn("Item", explode(col("ItemArray"))) .drop("ItemArray", "Items")
示例输出
原rankedDF中的第一行:
| TimePeriod | TPStartDate | TPEndDate | TXN_HEADER_ID | Items |
|---|---|---|---|---|
| 1 | 2017-03-01 | 2017-05-30 | TxnHeader1 | Womens Socks, Men Pants |
拆分后会变成两行:
| TimePeriod | TPStartDate | TPEndDate | TXN_HEADER_ID | Item |
|---|---|---|---|---|
| 1 | 2017-03-01 | 2017-05-30 | TxnHeader1 | Womens Socks |
| 1 | 2017-03-01 | 2017-05-30 | TxnHeader1 | Men Pants |
内容的提问来源于stack exchange,提问作者LeeFernan
相关产品推荐
相关产品推荐

