Question around joins. I generally see all my join...
# general
g
Question around joins. I generally see all my joins use shuffle hash join, which is different from what I see many times with Spark that uses sort merge by default I believe. Curious how Daft chooses the join strategy and how that differs from how spark chooses the strategy leading to differences in how both engines decide to join the same data.
Context is mostly around migrating code and tuning memory/cpu of jobs during the migration and what role this switch in join strategy could play during the migration (e.g. most of my spark jobs uses relatively low head node memory, but I have seen a need for higher memory in daft jobs).
k
Hi Garrett, in general the advantage of sort-merge join is to either deal with skew (e.g. if a column has a lot of the same value) or when the inputs are already sorted/need to be sorted.
🙏 1
However the join strategies are not fully comparable between the two engines since the implementations are different. The likely reason why our join uses larger head memory is because of how we use the Ray object store. We're aware that this is an issue and are actively working on improving the memory stability of joins. My guess is you'll likely see a pretty decent improvement in head node memory in the upcoming weeks/months