【问题标题】:How do I Window until a criteria is met within my Window?在我的窗口中满足某个条件之前,我如何窗口化?
【发布时间】:2021-05-29 09:59:13
【问题描述】:

我有一组包含 3 列感兴趣的数据。第一列是代表月份的日期。第二个是包含该月一些起始金额的列。第三个是表示该月金额减少的列。我有一段时间内每个月的多行数据。

例如,我们可能会得到一个日期为 2020-01-01,起始金额为 5MM,减少金额为 2MM。这意味着我们预计月底会有 3MM 的剩余金额。

我需要计算在接下来的几个月中消耗这个起始数量需要多长时间。

给定上面的例子,如果我们从 5MM 开始,那一个月消耗了 2MM,我们还剩下 3MM。如果下个月 2020-02-01 消耗 1.5MM,我们还剩下 1.5MM。如果下个月 2020-03-01 消耗了 2MM,我们还剩下 -0.5MM,我们在 2020-03-01 月份完成了消耗量。 2020-03-01 的这个结果是我希望得到的。

我怎样才能得到这个值?

我想在 DataFrame 的每一行中获得一个结果,并且我需要对 DataFrame 的其余部分执行聚合以查看历史行。因此,我假设我需要使用 Window 来计算这个值。但是,我无法弄清楚如何正确设置实际的 Window。

我在 Window 中的函数是采用称为“opening_amount”的起始量并在 Window 上减去称为“consume_amount”的燃尽量。例如,

  def followingWindowSpec: WindowSpec = Window.partitionBy(
    partitionCols : _*
  )
    .orderBy(orderByCols: _*)
    .rangeBetween(0, Window.unboundedFollowing)

  val burndownCompleteDateExpr = min(
    when(
      col("opening_amount")
        - sum(col("consume_amount"))
        .over(followingWindowSpec)
        <= lit(0),
      col("fiscal_dt")
    )
  )
    .over(followingWindowSpec)

我相信我需要使用一个从当前行开始并向前看的窗口。我尝试过使用 Window(0, x) ,其中 x 是某个值。

当我将两个 Windows 都设置为 UnboundedFollowing 时,我会得到每个月的财务数据。

当我将 consume_amount 总和窗口设置为使用 UnboundedPreceding 时,我得到了第一个月正确结果的正确结果,因为它具有 NULL 前面的值,但接下来的几个月要么返回同一个月(前几个月)或他们自己的月份(最初几个月之后)。

如果您能给我指点如何正确地进行窗口化,或者如果我在错误的树上吠叫正确的方法是什么,我将不胜感激。

样本数据:

+----+----------+-------+--------+-------------+
|item|      date|opening|consumed|out_of_supply|
+----+----------+-------+--------+-------------+
| 101|2020-01-01|   3200|    2000|   2020-02-01|
| 101|2020-02-01|   4600|    1500|         null|
| 101|2020-03-01|   1500|    1300|   2020-04-01|
| 101|2020-04-01|   4000|     500|         null|
| 220|2020-01-01|   3400|    2000|   2020-02-01|
| 220|2020-02-01|   1600|    3000|   2020-02-01|
| 220|2020-03-01|    310|    1000|   2020-03-01|
| 220|2020-04-01|    680|     500|         null|
+----+----------+-------+--------+-------------+

对于每一行,我将行中消耗的值向前n行相加,一次增加n个,看何时开始值被完全消耗。

比如2020-01-01的101项,开3200,第一个月消耗2000,月底1200。这1200被2020-02-01的1500完全消耗,所以 2020-02-01 是我正在寻找的月份。

对于 2020-02-01 中的第 101 项,它永远不会被完全消耗,因此我将默认返回 null。

【问题讨论】:

  • 如果每一行都有“开仓金额”和“消费金额”,为什么不简单地做一个groupBy()来聚合min(when($"opening_amt" &lt;= $"consumed_amt", $"fiscal_dt"))呢?
  • 问题是opening_amt在每一行都是现成的,但是consumed_amt需要累积,所以需要看其他行。
  • 鉴于以下 3 行 ("2020-01-01", 5000, 2000), ("2020-02-01", 3000, 1500), ("2020-03-01", 1500, 2000),第 3 行本身是否有足够的信息导致 true 处于 when() 条件中?
  • 第三行会,是的。前两个不会。
  • 前两行将导致null 因此被排除在外。因此,df.groupBy("item").agg(min(when($"opening_amt" &lt;= $"consumed_amt", $"fiscal_dt")).as("outofstock_date") 之类的内容应该可以满足您的需求。

标签: scala apache-spark


【解决方案1】:

我建议使用Window/rowsBetween() 为每一行组合以下行中的consumed 数据列表,然后由UDF 处理以捕获供应不足的date:

val df = Seq(
  (101, "2020-01-01", 3200, 2000),
  (101, "2020-02-01", 4600, 1500),
  (101, "2020-03-01", 1500, 1300),
  (101, "2020-04-01", 4000,  500),
  (220, "2020-01-01", 3400, 2000),
  (220, "2020-02-01", 1600, 3000),
  (220, "2020-03-01",  310, 1000),
  (220, "2020-04-01",  680,  500)
).toDF("item", "date", "opening", "consumed")

import org.apache.spark.sql.Row
import org.apache.spark.sql.expressions.Window
val win1 = Window.partitionBy("item").orderBy("date").
  rowsBetween(0, Window.unboundedFollowing)

val outOfStockDate = udf { (opening: Int, list: Seq[Row]) =>
  @scala.annotation.tailrec
  def loop(ls: List[(String, Int)], dt: String, acc: Int): String = ls match {
    case Nil =>
      null
    case head :: tail =>
      val accNew = acc + head._2
      if (opening <= accNew) head._1 else loop(tail, head._1, accNew)
  }
  loop(list.map{ case Row(d: String, c: Int) => (d, c) }.toList, null, 0)
}

df.
  withColumn("consumed_list", collect_list(struct($"date", $"consumed")).over(win1)).
  withColumn("out_of_supply", outOfStockDate($"opening", $"consumed_list")).
  drop("consumed_list").
  show
// +----+----------+-------+--------+-------------+
// |item|      date|opening|consumed|out_of_supply|
// +----+----------+-------+--------+-------------+
// | 101|2020-01-01|   3200|    2000|   2020-02-01|
// | 101|2020-02-01|   4600|    1500|         null|
// | 101|2020-03-01|   1500|    1300|   2020-04-01|
// | 101|2020-04-01|   4000|     500|         null|
// | 220|2020-01-01|   3400|    2000|   2020-02-01|
// | 220|2020-02-01|   1600|    3000|   2020-02-01|
// | 220|2020-03-01|    310|    1000|   2020-03-01|
// | 220|2020-04-01|    680|     500|         null|
// +----+----------+-------+--------+-------------+

【讨论】:

  • 感谢您的回答。实际用例中的一个混淆因素是每一行都有自己的开盘值,并且需要获取自己的日期。 IE。项目 101 在 1 月份的期初库存可能为 5000,但在 2 月份的期初库存为 3000。这也意味着结果不能按项目,而必须按项目和日期 - 这使它看起来更加复杂。
  • 不确定我是否理解您所说的“每一行都有自己的开盘价并且需要有自己的日期”是什么意思。也许您可以通过编辑您的问题来详细说明,以根据一些示例数据集添加所需的输出。
  • 抱歉耽搁了。我用另一种不太优雅的方法让它工作,并且有其他工作重点。我添加了示例数据来向您展示我的想法。我一直在空闲时间自己研究它,但尚未解决它。
  • 因此,根据您最新的样本数据,每行中的“开放”数量受外部变化的影响,无法从之前的“开放”/“消耗”数据中推断出来。请查看我修改后的解决方案。
  • 是的,开仓金额来自另一个系统,并且每个月都在变化。 IE。如果一个月内开仓是 3000,我们收到了 2000 的交付和消耗 1000,那么下个月的开仓是 4000。在他们向我们提供数据之前,我们不知道实际的交付和消耗是多少。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-09-09
  • 2020-05-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多