Spark filter dataframe on a condition based in the value on other rows

Viewed 30

I have a dataframe with 2 column start and end, I want to find all masters in the dataframe.

A master is defined like this: we order the dataframe by start ascending and end descending, the first row in this order is the first master, the next master is the row with a start value strictly bigger than the current master end, the same thing for this new master and so on ..

Example:

val columns = Seq("start", "end")
val data = Seq(
  (1, 5),
  (2, 7),
  (6, 9),
  (6, 9),
  (7, 12),
  (9, 14),
  (13, 15),
  (20, 30),
  (27, 29)
)
val df = spark.sparkContext.parallelize(data).toDF(columns: _*)
val w = Window.orderBy(col("start").asc, col("end").desc)
df.withColumn("rank", dense_rank().over(w)).show()

+-----+---+----+
|start|end|rank|
+-----+---+----+
|    1|  5|   1|
|    2|  7|   2|
|    6|  9|   3|
|    6|  9|   3|
|    7| 12|   4|
|    9| 14|   5|
|   13| 15|   6|
|   20| 30|   7|
|   27| 29|   8|
+-----+---+----+

The expected results is:

+-----+---+
|start|end|
+-----+---+
|    1|  5|
|    6|  9|
|   13| 15|
|   20| 30|
+-----+---+

My solution is to get the first row (rank = 1) get it's end value and inside a while loop get the next start that's superior to the current end so the current master will be that row and so on till no master is found.

Even though this work but it's obviously not good for large data, is there a way to do it only with spark functions (without loops) ?

0 Answers
Related