【发布时间】: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" <= $"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" <= $"consumed_amt", $"fiscal_dt")).as("outofstock_date")之类的内容应该可以满足您的需求。
标签: scala apache-spark