【问题标题】:Merging rows in a dataset合并数据集中的行
【发布时间】: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


    【解决方案1】:

    您提出的解决方案几乎是正确的,只是需要一些调整

    首先,由于您的方法 Row.merge 返回的是行列表而不是行,因此您应该使用 flatMapGroups 将列表分解为数据集中的不同记录:

    df
      .as[Row]
      .groupByKey(_.id)
      .flatMapGroups { case (k, iter) => Row.merge(iter) }
    

    接下来,让我们深入了解您的 Row.merge 方法。

    您创建一个空的bufferRow,在循环的第一次迭代中使用if (listBuffer.isEmpty) bufferRow = row 语句填充iter。但是,此语句中的条件对于所有迭代都是正确的,这就是您的输出仅包含每个组的最新行的原因。所以这个说法应该去掉。要初始化您的bufferRow,您只需调用iter.next()

    ... = new ListBuffer[Row]
    var bufferRow = iter.next()
    for(row <- iter) { ...
    

    由于迭代器来自 groupBy,它至少包含一个元素,因此第一次调用 iter.next() 不会抛出异常。并且由于对.next() 方法的调用删除了迭代器的第一个元素,之后的循环不会重新处理这个第一个元素。

    接下来,return 语句之前的方法的最后一条语句是if (listBuffer.isEmpty) listBuffer += bufferRow。这个语句不应该有条件。

    确实,在您的循环中,您填充bufferRow,然后仅当当前处理的行将“sw”字段设置为“on”时才将其添加到listBuffer。这个当前处理的行变成了新的bufferRow。这意味着最后一个bufferRow 永远不会保存在listBuffer 中,除非listBuffer 为空。所以merge 方法的最后几行应该是:

    ...
        bufferRow = row
      }
    }
    listBuffer += bufferRow
    

    我们现在有了完整的merge 方法:

    def merge(iter: Iterator[Row]): ListBuffer[Row] = {
      val listBuffer = new ListBuffer[Row]
      var bufferRow = iter.next()
      for (row <- iter) {
        if (row.sw == "off") bufferRow.end = row.end
        else if (row.sw == "on" ) {
          listBuffer += bufferRow
          bufferRow = row
        }
      }
      listBuffer += bufferRow
    }
    

    运行此代码会得到以下结果,按 id 和 start 列重新排序:

    +---+-----+---+---+
    |id |start|end|sw |
    +---+-----+---+---+
    |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|
    +---+-----+---+---+
    

    最后注意:如果在分区数据集上运行此代码,则应注意迭代器排序,我不确定 Spark 的 groupBy 方法是否保持按迭代器分组的行排序。也许在迭代之前用 .toSeq.sortBy(...) 重新排序迭代器更明智。

    【讨论】:

      猜你喜欢
      • 2020-12-19
      • 2021-10-10
      • 1970-01-01
      • 1970-01-01
      • 2020-12-06
      • 1970-01-01
      • 2022-11-11
      相关资源
      最近更新 更多