【发布时间】:2019-10-25 13:12:11
【问题描述】:
我有一个 spark 数据框,其中有几个列,如 tin、year、date_begin、date_end、continuous_data
tin year continuous_data
a1 2017 0
a1 2017 1
a1 2017 0
a1 2017 1
a1 2017 1
a1 2017 0
a1 2017 1
a1 2017 1
a1 2017 1
a1 2017 0
a1 2017 1
同样,我还有 2 个日期时间格式为 (yyyy-mm-dd HH:mm:ss) 的列。
我需要访问 'continuous_data' 列的每一行,例如 x(i+1) 和 x(i-1)。就我而言,就像
continuous_data(i) - 当前行值
Continuous_data(i-1) - 上一行值
Continuous_data(i+1) - 下一行值
这样我的需求如下所示
tin year continuous_data prev_data next_data
a1 2017 0 null 1
a1 2017 1 0 0
a1 2017 0 1 1
a1 2017 1 0 1
a1 2017 1 1 0
a1 2017 0 1 1
a1 2017 1 0 1
a1 2017 1 1 1
a1 2017 1 1 0
a1 2017 0 1 1
a1 2017 1 0 null
我需要在纯 Scala 中解决它,而不是使用 spark 函数,我使用窗口函数来实现它,由于某些原因不需要。
我试图从过去几天解决这个问题,但我还不能解决它。有人可以帮我解决这个问题。
【问题讨论】:
-
只是为了好奇,为什么不能使用窗口函数呢?常规 scala 集合具有
sliding函数,但数据集没有,最简单的解决方法是使用窗口函数。 -
@KrzysztofAtłasik,是的,这是最简单的方法,我也这样做了,虽然没问题,但我必须用纯 scala 脚本实现相同的方法,特别是 x(i+1) , x(i-1)
-
您可以使用 Spark Udafs 来完成。如果你支持这个选项,我可以给你看一些代码
-
@EmiCareOfCell44 是的,您能否分享您对此的想法。但我主要要做的事情是 x(i+1), x(i-1)。如果使用 udaf 可以实现,我会尝试实现它,谢谢
-
其背后的想法是 udaf 可以处理您使用 Scala 集合处理数据的聚合。它意味着从数据帧格式反序列化,但您可以以通用方式处理您的窗口,应用 f: Seq[(A, C)] => Seq[(A, C)] 之类的函数。我会放代码
标签: scala apache-spark-sql rdd