Given the dimensions of the datasets I am currently using, I started working in Databricks with PySpark. After a few weeks I still struggle to fully understand what is going on under the hood.
I have a dataset of about 40 millions rows and I apply this function that dynamically compute some aggregation in a rolling window:
def freshness(df):
days = lambda x: x*60*60*24
w = Window.partitionBy('csecid', 'date').orderBy('date')
w1 = Window.partitionBy('csecid').orderBy(F.col('date').cast('timestamp').cast('long')).rangeBetween(-days(100), 0)
w2 = Window.partitionBy('csecid').orderBy('date')
w3 = Window.partitionBy('csecid', 'id', 'date').orderBy('date')
w4 = Window.partitionBy('csecid', 'id')
w5 = Window.partitionBy('csecid', 'id').orderBy(F.col('date').desc())
df = df.withColumn('dateid', F.row_number().over(w))
df1 = df.withColumn('flag', F.collect_list('date').over(w1))
df1 = df1.withColumn('id', F.row_number().over(w2)).select('csecid', 'flag', 'id')
df1 = df1.withColumn('date', F.explode('flag')).drop('flag')
df1 = df1.withColumn('dateid', F.row_number().over(w3))
df2 = df1.join(df, on=['csecid', 'date', 'dateid'], how='left')
df3 = (df2
.withColumn('analyst_fresh', F.floor(F.approx_count_distinct('analystid', 0.005).over(w4)/3))
.orderBy('csecid', 'id', 'date', F.col('perenddate').desc())
.groupBy('csecid', 'id', 'analystid')
.agg(F.last('date', True).alias('date'), F.last('analyst_fresh', True).alias('analyst_fresh'))
.where(F.col('analystid').isNotNull())
.orderBy('csecid', 'id', 'date')
.withColumn('id2', F.row_number().over(w5))
.withColumn('freshness', F.when(F.col('id2')<=F.col('analyst_fresh'), 1).otherwise(0))
.drop('analyst_fresh', 'id2', 'analystid')
)
df_fill = _get_fill_dates_df(df3, 'date', ['id', 'csecid'])
df3 = df3.join(df_fill, on=['csecid', 'date', 'id', 'freshness'], how='outer')
df3 = df3.groupBy('csecid', 'id', 'date').agg(F.max('freshness').alias('freshness'))
df3 = df2.join(df3, on=['csecid', 'id', 'date'], how='left').fillna(0, subset=['freshness'])
df3 = df3.withColumn('fresh_revision', (F.abs(F.col('revisions_improved'))+F.col('freshness'))*F.signum('revisions_improved'))
df4 = (df3
.orderBy('csecid', 'id', 'date', F.col('perenddate').desc())
.groupBy('csecid', 'id', 'analystid')
.agg(F.last('date').alias('date'), F.last('fresh_revision', True).alias('fresh_revision'))
.orderBy('csecid', 'id', 'date')
.groupBy('csecid', 'id').agg(F.last('date').alias('date'),
F.sum('fresh_revision').alias('fresh_revision'), F.sum(F.abs('fresh_revision')).alias('n_revisions'))
.withColumn('revision_index_improved', F.col('fresh_revision') / F.col('n_revisions'))
.groupBy('csecid', 'date')
.agg(F.first('revision_index_improved').alias('revision_index_improved'))
)
df5 = df.join(df4, on=['csecid', 'date'], how='left').orderBy('csecid', 'date')
return df5
weight_list = ['leader', 'truecall']
for c in weight_list:
df = df.withColumn('revisions_improved', (F.abs(F.col('revisions_improved'))+F.col(c))*F.col('revisions'))
df = freshness(df)
This piece of code run in about 2 hours, and I noticed that most of the computational time is taken from one job that uses a single executor (out of 16 available).
I read that it is possible to overcome this problem by repatitioning the dataframe with .repartition(). My questions are the following: how can I find which piece of code the abovementioned job is referred to? Where should I repartition my dataframe? Is that a correct solution?

