Is daft also implementing exoshuffle currently? Or some other shuffling mechanism?
j
jay
10/14/2024, 8:29 AM
Exoshuffle is a paper that talks about using Ray to implement various shuffling algorithms
We’re working on bringing the push-based shuffle in the paper to Daft. We have a prototype available in a PR already which has been really promising, but more details soon by end of month!
Are you running into shuffling limits atm?
k
Kyle
10/14/2024, 8:32 AM
oh nice! was just hoping to get faster joins haha
Kyle
10/14/2024, 8:34 AM
Came across it, thought it was totally in your wheelhouse, and also thought you'd probably be able to get better results than their ray data example 😄
I was able to push this quite easily to a 3000x3000 partition shuffle on Ray, which previously would make our naive shuffle choke pretty hard (since it would create 9M intermediate objects in Ray)
I think it scales well even beyond that