Choosing a partition size
So how do I choose the argument for partitionBy()? How many partitions is the right partition size? If you have too few partitions, it won't take full advantage of your cluster, it can't spread it out effectively. On the other hand, if you have too many you end up shuffling data around and breaking things up into chunks that are too small; there's some overhead associated with running an individual executor job, so you don't want too many executors either. You want at least as many partitions as you have cores in your cluster or as many will fit in your available memory. A hundred is usually a reasonable place to start. Let's say you have five or ten computers in your cluster-that will split it up into a reasonable ...
Become an O’Reilly member and get unlimited access to this title plus top books and audiobooks from O’Reilly and nearly 200 top publishers, thousands of courses curated by job role, 150+ live events each month,
and much more.
Read now
Unlock full access