Spark stage taking too long - 2 executors doing "all" the work

Viewed 915

I've been trying to figure this out for the past day, but have not been successful.

Problem I am facing

I'm reading a parquet file that is about 2GB big. The initial read is 14 partitions, then eventually gets split into 200 partitions. I perform seemingly simple SQL query that runs for 25+ mins runtime, about 22 mins is spent on a single stage. Looking in the Spark UI, I see that all computation is eventually pushed to about 2 to 4 executors, with lots of shuffling. I don't know what is going on. Please I would appreciate any help.

Setup

  • Spark environment - Databricks
  • Cluster mode - Standard
  • Databricks Runtime Version - 6.4 ML (includes Apache Spark 2.4.5, Scala 2.11)
  • Cloud - Azure
  • Worker Type - 56 GB, 16 cores per machine. Minimum 2 machines
  • Driver Type - 112 GB, 16 cores

Notebook

Cell 1: Helper functions


load_data = function(path, type) {
  input_df = read.df(path, type)

  input_df = withColumn(input_df, "dummy_col", 1L)
  createOrReplaceTempView(input_df, "__current_exp_data")

  ## Helper function to run query, then save as table
  transformation_helper = function(sql_query, destination_table) {
    createOrReplaceTempView(sql(sql_query), destination_table)
  }

  ## Transformation 0: Calculate max date, used for calculations later on
  transformation_helper(
    "SELECT 1L AS dummy_col, MAX(Date) max_date FROM __current_exp_data",
    destination_table = "__max_date"
  )

  ## Transformation 1: Make initial column calculations
  transformation_helper(
    "
    SELECT
        cId                                                          AS cId
      , date_format(Date, 'yyyy-MM-dd')                              AS Date
      , date_format(DateEntered, 'yyyy-MM-dd')                       AS DateEntered
      , eId

      , (CASE WHEN isnan(tSec) OR isnull(tSec) THEN 0 ELSE tSec END) AS tSec
      , (CASE WHEN isnan(eSec) OR isnull(eSec) THEN 0 ELSE eSec END) AS eSec

      , approx_count_distinct(eId) OVER (PARTITION BY cId)           AS dc_eId
      , COUNT(*)                   OVER (PARTITION BY cId, Date)     AS num_rec

      , datediff(Date, DateEntered)                                  AS analysis_day
      , datediff(max_date, DateEntered)                              AS total_avail_days
    FROM __current_exp_data
    CROSS JOIN __max_date ON __main_data.dummy_col = __max_date.dummy_col
  ",
    destination_table = "current_exp_data_raw"
  )

  ## Transformation 2: Drop row if Date is not valid
  transformation_helper(
    "
    SELECT 
        cId
      , Date
      , DateEntered
      , eId
      , tSec
      , eSec

      , analysis_day
      , total_avail_days
      , CASE WHEN analysis_day     == 0 THEN 0    ELSE floor((analysis_day - 1) / 7)   END AS week
      , CASE WHEN total_avail_days  < 7 THEN NULL ELSE floor(total_avail_days / 7) - 1 END AS avail_week
    FROM current_exp_data_raw
    WHERE 
        isnotnull(Date) AND 
        NOT isnan(Date) AND 
        Date >= DateEntered AND
        dc_eId == 1 AND
        num_rec == 1
  ",
    destination_table = "main_data"
  )
  
  cacheTable("main_data_raw")
  cacheTable("main_data")
}

spark_sql_as_data_table = function(query) {
  data.table(collect(sql(query)))
}

get_distinct_weeks = function() {
  spark_sql_as_data_table("SELECT week FROM current_exp_data GROUP BY week")
}

Cell 2: Call helper function that triggers the long running task

library(data.table)
library(SparkR)

spark = sparkR.session(sparkConfig = list())

load_data_pq("/mnt/public-dir/file_0000000.parquet")

set.seed(1234)

get_distinct_weeks()

Long running stage DAG

Long running stage DAG

Stats about long running stage

Executor stats

Logs

I trimmed it down, and show only entries that appeared multiple times below


BlockManager: Found block rdd_22_113 locally

CoarseGrainedExecutorBackend: Got assigned task 812

ExternalAppendOnlyUnsafeRowArray: Reached spill threshold of 4096 rows, switching to org.apache.spark.util.collection.unsafe.sort.UnsafeExternalSorter

InMemoryTableScanExec: Predicate (dc_eId#61L = 1) generates partition filter: ((dc_eId.lowerBound#622L <= 1) && (1 <= dc_eId.upperBound#621L))

InMemoryTableScanExec: Predicate (num_rec#62L = 1) generates partition filter: ((num_rec.lowerBound#627L <= 1) && (1 <= num_rec.upperBound#626L))

InMemoryTableScanExec: Predicate isnotnull(Date#57) generates partition filter: ((Date.count#599 - Date.nullCount#598) > 0)

InMemoryTableScanExec: Predicate isnotnull(DateEntered#58) generates partition filter: ((DateEntered.count#604 - DateEntered.nullCount#603) > 0)

MemoryStore: Block rdd_17_104 stored as values in memory (estimated size <VERY SMALL NUMBER &lt; 10> MB, free 10.0 GB)

ShuffleBlockFetcherIterator: Getting 200 non-empty blocks including 176 local blocks and 24 remote blocks

ShuffleBlockFetcherIterator: Started 4 remote fetches in 1 ms

UnsafeExternalSorter: Thread 254 spilling sort data of <Between 1 and 3 GB> to disk (3  times so far)

0 Answers
Related