Skip to content

mau-foo/spark-challenges

 
 

Folders and files

NameName
Last commit message
Last commit date

Latest commit

 

History

10 Commits
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Data skew problem

Problem description

Running joins of skewed data sets we can face stragglers during execution which slow down whole app. Detail description in the blog post

Solution

Such slow tasks are mostly results of non-uniform data distribution between workers of spark app. The possible solution is to redistribute large data set between available workes in uniform manner. And in the same time, broadcast smaller data set to all workers pre-sorted by column for consolidation.

Solution of skew problem

  • uniform repartition for large data set (partition number depends from available resources)
  • broadcast for sorted smaller data set (could be shrinked to fewer attributes)

Data set could be repartitioned with

  • DataFrame default repartitioner
  • DataFrame repartitioner with special expression based on data columns
  • RDD custom repartitioner

Repartition strategy depends of data nature and domain specifics. Are preferred DataFrame/DataSets API.

Additionally, we can tweak solution with executors and cores number for better cluster resources utilisation.

Run app

step 1. build app assembly

$ sbt clean compile assembly

step 2. copy assembly jar to cluster step 3. run spark job on cluster with command

spark-submit --class DataSkew --master yarn --deploy-mode client sparkchallenges-assembly-0.1.jar <users-dataset-path> <deps-dataset-path>

or with cluster resources tweak (3 executors with 3 cores each)

spark-submit --class DataSkew --master yarn --deploy-mode client --num-executors 3 --driver-memory 1G --executor-memory 3G --executor-cores 3 sparkchallenges-assembly-0.1.jar <users-dataset-path> <deps-dataset-path>

References

Have fun! Team @DataEngi

About

data engineering challenges and fun

Resources

License

Stars

Watchers

Forks

Releases

No releases published

Packages

No packages published

Languages

  • Scala 100.0%