Spark Shuffle Write Size 7500 times the Input Size

Viewed 103

I'm dealing with a very odd situation in one of our spark clusters.

We are testing a spark application which runs in client mode and use dynamic allocation for the executors.

While testing the application, QA decided to manually kill it (kill command, not Spark UI) and restart it. Once the application started processing data again it got stuck forever in one of the stages. After the problem happened the first time, the behavior of the application was always the same, getting stuck always at the same place (collectAsList(). At this point the data is expected to be very small, enough to fit in the drivers memory).

Although I was able to reproduce the issue in the problematic environment, I couldn't see this behavior on another environment. Looking for clues in the Spark UI I noticed that although it loaded 15k records from kafka (which is nothing), the input it was processing was showing only 2 records and somehow, the size of the shuffle write was a staggering 6.3GB for 17298252119 records. This stage was running using 4 partitions only (we are using hash partition with the default 200 setting for shuffle).

Shuffle Write Size

The only theory I have so far is that killing the process somehow corrupted some spark or shuffle service tracking info and instead of processing just the data it received, the shuffle service is trying to restore old records.

One additional piece of info that seems to be linked to shuffle issues, is that the application in this case processes the dataset multiple times and then joining all the resulting data frames (in this case there are 4 datasets joined).

Any help is welcome. Although I'm leaning to conclude the environment got corrupted, I cannot prove it and since spark is a complex beast, this issue might happen again.

Thank you.

0 Answers
Related