Enhancing performance Pyspark RDD

Viewed 24

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).

enter image description here

enter image description here

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?

0 Answers
Related