【问题标题】:Filtering a spark DataFrame down to updates only将 spark DataFrame 过滤为仅更新
【发布时间】: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


【解决方案1】:

您可以轻松使用lag。窗口

val window = Window.partitionBy($"id").orderBy($"date".asc)

import org.apache.spark.sql.functions.{coalesce, lag, lit}

val keep = coalesce(lag($"x", 1).over(window) =!= $"x", lit(true))

df.withColumn("keep", keep).where($"keep").drop("keep").show

// +--------+---+---+
// |    date| id|  x|
// +--------+---+---+
// |20150101|  1|  1|
// |20150105|  1|  2|
// |20150107|  1|  1|
// +--------+---+---+

【讨论】:

  • 这适用于我展示的示例,其中我有一个“数据”变量x,我想观察它的变化。但是如果我有多个列(x1x2,...,xn,比如说),我只知道来自Seq[String] 的那些列被称为(以及有多少)。这种情况下的问题是lag 不是可变参数。我是否应该制作不同版本的"keep"Seq[Column],然后通过列的顺序折叠线程df 并过滤它们?
  • 你需要像val keep = cols.map(x => coalesce(lag(x, 1).over(window) =!= x, lit(true))).reduce(_ | _)这样的东西
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-03-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多