【发布时间】:2019-01-31 17:18:05
【问题描述】:
我有一个时间序列数据集,它按 id 分区,并按时间戳排序。示例:
ID Timestamp Feature
"XSC" 1986-05-21 44.7530
"XSC" 1986-05-22 44.7530
"XSC" 1986-05-23 23.5678
"TM" 1982-03-08 22.2734
"TM" 1982-03-09 22.1941
"TM" 1982-03-10 22.0847
"TM" 1982-03-11 22.1741
"TM" 1982-03-12 22.1840
"TM" 1982-03-15 22.1344
我有一些我需要计算的自定义逻辑,它应该在每个窗口、每个分区内完成。 我知道 Spark 对窗口函数有丰富的支持,我正在尝试将其用于此目的。
我的逻辑需要当前窗口/分区中的元素总数,作为标量。我需要它来做一些特定的计算(基本上,一个 for 循环最多)。
我尝试添加一个计数列,方法是
val window = Window.partitionBy("id").orderBy("timestamp")
frame = frame.withColumn("my_cnt", count(column).over(window))
我需要做类似的事情:
var i = 1
var y = col("Feature")
var result = y
while (i < /* total number of records within each partition goes here */) {
result = result + lit(1) * lag(y, i).over(window) + /* complex computation */
i = i + 1
}
dataFrame.withColumn("Computed_Value", result)
如何将每个分区中的记录总数作为标量值?我还添加了计数“my_cnt”值,它添加了分区的总值,但在我的情况下似乎无法使用它。
【问题讨论】:
-
你能展示一些示例输入和预期的输出吗?不清楚您要做什么?
-
添加了输入和一些示例代码。
-
还是不清楚。根据您的输入清楚地提供预期的输出。
-
也许你需要一个聚合窗口函数,比如blog.nuvola-tech.com/2017/10/…
-
@1pluszara - 这里的问题不在于输出应该是什么。没关系。重要的是,我如何访问当前窗口/分区中的元素总数。粘贴在那里的代码只是一些逻辑来查看我需要总计数的位置,以及我需要它的格式(作为实际值,而不是列)
标签: apache-spark apache-spark-sql