【发布时间】:2016-04-02 15:13:50
【问题描述】:
我有一个数据框(火花):
id value
3 0
3 1
3 0
4 1
4 0
4 0
我想创建一个新的数据框:
3 0
3 1
4 1
需要为每个 id 删除 1(value) 之后的所有行。我尝试使用 spark dateframe(Scala) 中的窗口函数。但无法找到解决方案。似乎我走错了方向。
我正在寻找 Scala 中的解决方案。谢谢
使用 monotonically_increasing_id 输出
scala> val data = Seq((3,0),(3,1),(3,0),(4,1),(4,0),(4,0)).toDF("id", "value")
data: org.apache.spark.sql.DataFrame = [id: int, value: int]
scala> val minIdx = dataWithIndex.filter($"value" === 1).groupBy($"id").agg(min($"idx")).toDF("r_id", "min_idx")
minIdx: org.apache.spark.sql.DataFrame = [r_id: int, min_idx: bigint]
scala> dataWithIndex.join(minIdx,($"r_id" === $"id") && ($"idx" <= $"min_idx")).select($"id", $"value").show
+---+-----+
| id|value|
+---+-----+
| 3| 0|
| 3| 1|
| 4| 1|
+---+-----+
如果我们在原始数据框中进行排序转换,该解决方案将不起作用。那个时候monotonically_increasing_id()是基于原始DF而不是排序DF生成的。我之前错过了这个要求。
欢迎所有建议。
【问题讨论】:
-
到目前为止你尝试了什么?
-
@eliasah 我根据stackoverflow.com/questions/32148208/… 的答案尝试了一些实验。但到目前为止没有成功
-
你的 DF 排序了吗?
-
@TheArchetypalPaul 是的,它已排序
-
因为你每次都调用
show。在我下面的代码中,评估是懒惰的——原始的val dataWithIndex仅在我的最终show被调用时才被评估。但是你每次都打电话show,迫使重新评估。停止调用show,或创建dataWithIndex后立即调用cache
标签: scala apache-spark dataframe apache-spark-sql