How to efficiently split a dataframe in Spark based on a condition?

Viewed 60

I have a situtation like that with this Spark dataframe:

id value
1 0
1 3
2 4
1 0
2 2
3 0
4 1

Now what I want to obtain is to efficiently split this single dataframe in 3 different one such that each dataframe extracted from the original one is between two 0 in the "value" column (with the first zero indicating the beginning of each dataframe) using Apache Spark, so that I would obtain this as result:

Dataframe 1 (rows from first 0 value to the last value before the next 0):

id value
1 0
1 3
2 4

Dataframe 2 (rows from the second zero value to the last value before the 3rd zero):

id value
1 0
2 2

Dataframe 3:

id value
3 0
4 1
1 Answers

While as samkart said it is not efficient/easy way to break data on basis of order of rows still if you are using spark v3.2+ you can leverage pandas on pyspark to do it in spark way like below

import pyspark.pandas as ps
from pyspark.sql import functions as F
from pyspark.sql import Window
pdf=ps.read_csv("/FileStore/tmp4/pand.txt")
sdf = pdf.to_spark(index_col='index')
sdf=sdf.withColumn("run",F.sum(F.when(F.col("value")==0,1).otherwise(0)).over(Window.orderBy("index")))
toval= sdf.agg(F.max(F.col("run"))).collect()[0][0]
for x in range (1,toval+1):
  globals()[f"sdf{x}"]=sdf.filter(F.col("run")==x).drop("index","run")

For above data it will create 3 dataframe sdf1,sdf2,sdf3 like below

sdf1.show()
sdf2.show()
sdf3.show()
#output
+---+-----+
| id|value|
+---+-----+
|  1|    0|
|  1|    3|
|  2|    4|
+---+-----+

+---+-----+
| id|value|
+---+-----+
|  1|    0|
|  2|    2|
+---+-----+

+---+-----+
| id|value|
+---+-----+
|  3|    0|
|  4|    1|
+---+-----+
Related