【发布时间】:2018-11-21 02:55:44
【问题描述】:
假设我在 Spark 中有一个 DataFrame,它由 id、date 和许多属性(例如 x、y、z)列组成。不幸的是,DataFrame 非常大。幸运的是,大多数记录都是“无变化”记录,其中 id、x、y、z 相同,只有日期发生变化。例如
| date | id | x |
| -------- | -- | - |
| 20150101 | 1 | 1 |
| 20150102 | 1 | 1 |
| 20150103 | 1 | 1 |
| 20150104 | 1 | 1 |
| 20150105 | 1 | 2 |
| 20150106 | 1 | 2 |
| 20150107 | 1 | 2 |
| 20150108 | 1 | 2 |
可以简化为
| date | id | x |
| -------- | -- | - |
| 20150101 | 1 | 1 |
| 20150105 | 1 | 2 |
我原本以为这个函数会做我想做的事
def filterToUpdates (df : DataFrame) = {
val colsData = df.column.filter(x => (x != "id" && x != "date"))
val window = Window.partitionBy(colsData).orderBy($"date".asc)
df.withColumn("row_num", row_number.over(window)).
select($"row_num" === 1).drop("row_num")
但是,如果我的数据列发生变化,然后再改回来,这会失败。
例如
| date | id | x |
| -------- | -- | - |
| 20150101 | 1 | 1 |
| 20150102 | 1 | 1 |
| 20150103 | 1 | 1 |
| 20150104 | 1 | 1 |
| 20150105 | 1 | 2 |
| 20150106 | 1 | 2 |
| 20150107 | 1 | 1 |
| 20150108 | 1 | 1 |
会变成
| date | id | x |
| -------- | -- | - |
| 20150101 | 1 | 1 |
| 20150105 | 1 | 2 |
而不是我想要的:
| date | id | x |
| -------- | -- | - |
| 20150101 | 1 | 1 |
| 20150105 | 1 | 2 |
| 20150107 | 1 | 1 |
在按顺序(按 id 分区并按日期排序)传递行记录的程序代码中,这不是一项艰巨的任务,但我只是看不出如何将其表述为 spark 计算。
注意:这与Spark SQL window function with complex condition 不同。我希望过滤掉与先前行不同的行,这是在该问题中构建became_inactive 列之后可以完成的另一件事。
【问题讨论】:
标签: apache-spark apache-spark-sql