Is daft also implementing exoshuffle currently? Or...
# general
k
Is daft also implementing exoshuffle currently? Or some other shuffling mechanism?
j
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
oh nice! was just hoping to get faster joins haha
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 😄
j
Yes 😛 here is the PR: https://github.com/Eventual-Inc/Daft/pull/2883 @Colin Ho is going to work on pushing this through after we get our streaming execution out
daft party 1
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
k
Super cool!! Can't wait! Haha