【发布时间】:2020-11-20 09:19:04
【问题描述】:
输入数据
val df = Seq(
("1", 1, 1, "on"),
("1", 2, 2, "off"),
("1", 2, 5, "off"),
("1", 5, 5, "on"),
("1", 5, 6, "off"),
("2", 1, 1, "off"),
("2", 1, 2, "off"),
("2", 2, 2, "on"),
("2", 3, 4, "off"),
("2", 5, 7, "off"),
("2", 8, 10, "on"),
("2", 11, 11, "on"),
("2", 11, 12, "off"),
("3", 1, 12, "off")
).toDF("id", "start", "end", "sw")
我正在尝试使用 groupBy 和 mapGroups 合并行。
期望的输出
1 1 5 on
1 5 6 on
2 1 2 off
2 2 7 on
2 8 10 on
2 11 12 on
3 1 12 off
逻辑如下。每个 off 行都合并到前一个 on 行中。如果第一个或唯一的值是关闭的,我会得到一个单独的行。从第一行开始,从最后一行结束。数据应按开始和结束排序。
这是我目前所拥有的
df
.as[Row]
.groupByKey(_.id)
.mapGroups{case(k, iter) => Row.merge(iter)}
我按 id 对数据进行分组,然后尝试迭代其他值。
case class Row(id:String, start:Int, var end:Int, sw:String)
object Row {
def merge(iter: Iterator[Row]): ListBuffer[Row] = {
val listBuffer = new ListBuffer[Row]
var bufferRow = Row("", 0, 0, "")
for(row <- iter){
if(listBuffer.isEmpty) bufferRow = row
else if(row.sw == "off") bufferRow.end = row.end
else if(row.sw == "on") {
listBuffer += bufferRow
bufferRow = row
}
}
if(listBuffer.isEmpty) listBuffer += bufferRow
listBuffer
}
}
我的输出
[WrappedArray([1,5,6,off])]
[WrappedArray([2,11,12,off])]
[WrappedArray([3,1,12,off])]
我已经使用窗口函数和累积和完成了类似的事情。在这里,我正在尝试学习一种新方法。
使用 spark 2.2、scala 2.11。
【问题讨论】:
标签: scala apache-spark