- Sources: Sem Sinchenko, graphframes-rs, discussion
- Summary: Sem Sinchenko's post, dated 2026-07-05 and resurfaced on Hacker News on 2026-07-31, implements Pregel-style bulk-synchronous graph algorithms on Apache DataFusion, offloading state to disk and relying on bulk scans rather than random access, which lets DataFusion handle spill, sort-merge joins, aggregation, and planning. PageRank on graph500-26, 32,804,978 nodes and 1,051,922,853 edges, ran under a 5 GB memory limit with a 4 GB DataFusion pool and took about 30 minutes for 15 iterations, matching ground truth to a 0.0001 tolerance. Weakly connected components on twitter_mpi, 52,579,682 nodes and 1,963,263,821 directed edges, ran under a 10 GB limit with an 8 GB pool, symmetrizing pushed the working set to 3,228,212,374 edges at peak, and the run completed in 22 forward iterations with component sizes matching the reference.
- Why it matters: The memory constraint is enforced rather than asserted, with the runs executed under
systemd-run using MemoryMax and MemorySwapMax=0 and both results checked against the Graphalytics ground truth, and the author names what does not work, including FairSpillPool deadlocks in extreme scenarios and no way found to make sort-merge join use pre-sorted data on disk.
send feedback on this story