I understand the concept of the broadcast optimization.
When one of the sides in the join have small data it's better to do the shuffle just for the small side.
but why isn't it possible to do this shuffle using only the executors? Why do we need to use the driver?
If each executor hold hash table to map the records between the executors I think it should work.
In the current implementation of spark broadcast - it collect the data to the driver and then shuffle it and the collect action to the driver is bottleneck that I would like to avoid.
Any ideas of how to achieve similar optimization without having the bottleneck of the driver memory?