Algorithms on billion-scale graph using 10GB RAM: I love DataFusion
Posted by speckx 3 days ago
Comments
Comment by chrisweekly 3 days ago
Impressive!
Comment by danbruc 3 days ago
Comment by eeks 2 days ago
[1] "Scalability! But at what COST?"
Comment by anon7725 3 days ago
You can get pretty far with sparse graphs, which are just arrays, in combination with memory mapping.
Comment by ssinchenko 2 days ago
Comment by nylonstrung 2 days ago
The extensibility is insane, you can create your own query language that compiles to logical plans.
Comment by ssinchenko 2 days ago
Comment by stuhood 18 hours ago
Comment by yadgire7 3 days ago
I am here to seek guidance from the community. I want to refresh my memory on knowledge graphs and algorithms for Big Data Mining and Processing.
I believe KG can solve problems on Agent attacks (LLM agency) in real-time - so want to build knowledge around the topic.
Interested to join any interest/ discussion groups if any. Thanks!
Comment by cpdomina 3 days ago
Comment by ssinchenko 2 days ago
Comment by adsharma 2 days ago
I forked networkit for exactly this reason. Columnar memory is much more efficient.
Comment by BigTTYGothGF 2 days ago
Comment by adsharma 3 days ago
https://github.com/Ladybug-Memory/icebug
Out of core with datafusion is the main innovation here in graphframes-rs. But it has only 2 algorithms so far.
Icebug and LadybugDB can be tightly integrated to efficiently move tables encoded as compressed sparse row (CSR) into arrow memory.
Jupyter notebooks available.
Comment by adsharma 3 days ago
Trade-off: datafusion allows you to do fine grained storage integration (spill to disk as a part of the algorithm).
The icebug/ladybug way is coarse grained. But it allows you to run cypher instead of writing datafusion operators.
Comment by lmeyerov 2 days ago
~10 years ago, we helped create apache arrow, helped create GPU data frames, and been running for the last decade the open source pygraphistry and now gfql cpu+gpu property graph engine for this. Likewise, Nvidia has been doing great with cuGraph (OSS) for GPU algs around this.
It is great you are finding success with this direction, but "shoulders of giants" merit credit - 100+ people.
Comment by adsharma 2 days ago
I didn't write the graph algorithms in networkit. Somebody else did.
But if you're looking for 100-200 graph algorithms running efficiently on columnar memory, I haven't found a working implementation that's permissively open source.
Happy to work with cugraph or anyone else who builds on top of Parquet/Arrow ecosystem.
Comment by lmeyerov 2 days ago
Comment by adsharma 2 days ago
https://arrow.apache.org/docs/r/authors.html
CUDA and the ecosystem around it is more complicated.
I'll use it. I agree that it advanced the state of the art and helped fund some of the truly OSS projects.
Don't feel the need to bring it up on a HN comment. Neither does the parent article by Sem.
In fact, it's in the vendor's interest to transparently route the algorithms in Icebug to the GPU in a compatible way instead of writing their own APIs.
Comment by adsharma 2 days ago
https://github.com/timlrx/graph-benchmarks
Not clear if the author or anyone else has an updated version of these benchmarks. Opened an icebug issue on the repo.
Comment by lmeyerov 2 days ago
we use igraph cpu / cugraph gpu in pygraphistry/gfql projects, and our experience is cugraph is ~10X+ on tiny cheap GPUs over igraph. looking at that table, where all seem same magnitude as igraph, I'd therefore expect much better than all the options listed. And when things get bigger and it merits a bigger GPU, even better. The team has been doing a great job over the last decade, and all free.
Comment by lmeyerov 2 days ago
can't reply on the below, but re:networkx, it was more of the reverse, they did a nice job of building a standalone embeddable table+graph arrow-friendly library over many years, and with networkx compatibility from the beginning. Later, they collaborated with the networkx team to eventually upstream it so networkx users can benefit more easily
as it is embeddable, you don't need to go through networkx to use it -- eg, we use it directly, as do various databases
and I'm not sure why you're saying people do not use networkx->cugraph because of performance, again, this is graph500 level performance that even graph database vendors now support because it is literally magnitudes faster than their cpu alternatives. We have had projects like court cases where the data science team switched to sitting on top of pygraphistry/gfql -> cudf/cugraph and using the GPU versions were the difference between hours and minutes, which for iterating over interactive analysis, is important. It's night and day going from CPU -> GPU, and a lot of butts were saved because of this.
We see the same thing with databricks+neo4j migrations.. most tools have their sweet spots, but also their comparative weak spots.
Comment by adsharma 2 days ago
fast on the GPU: nx-cugraph
slow on the CPU: networkx (no arrow, see linked graph-benchmark)
I'm looking to offer something that's best in class on both CPU and GPU, so people don't have to choose.Now that nvidia is becoming a major CPU vendor, they may like the idea too.
Comment by lmeyerov 2 days ago
The HPC lessons for graph algorithms is less obvious as the substrate the OP is working through is close but not quite what HPC folks figured out, so the trade-off of fast to write and maintainable (dataframe/db-based) vs at the achievable magnitude of performance is tricky. LLMs change the calculus here too IMO, so I've been thinking a lot about more NUMA-exposed ideas that before were relegated to PhD land.
Comment by ratmice 3 days ago
did datafusion gain some feature that they noted was missing in the previous article, or did something in their understanding click so they could overcome the previous issues?
Comment by ssinchenko 2 days ago
I tried using DataFusion as an in-memory tool, which was a mistake. If the graph fits in memory, Networkit, IGraph, etc. will almost always be faster. These tools cannot process anything bigger than the available memory.
So, I changed my approach. I wrote my own naive "disk checkpointer," offloading everything to disk and avoiding materialization. Although I was afraid that writing to and reading from the disk would be slow, it is surprisingly fast with DataFusion. The results are impressive: fast and out-of-core.
Sorry, this post is short and not very detailed. I did not expect it to be at the top of HN and receive so much attention.
Comment by adsharma 2 days ago
https://github.com/Ladybug-Memory/icebug-graphframes-compari...
I thought of datafusion work and the rust implementation as a way to address the delta seen vs the Spark/JVM GraphFrames implementation (author is a major contributor).
Looking forward to more such innovations, which will benefit the ecosystem as a whole. Why would anyone want to use a pure python graph algorithm package unless they're dealing with toy graphs?
Comment by lmeyerov 2 days ago
So then the question became pandas/polars/datafusion/duckdb/etc, must of which are rust/native. I'm curious myself why datafusion vs others, it's an interesting project :)
Comment by theLiminator 3 days ago
Comment by RobinL 2 days ago
Comment by ssinchenko 2 days ago
Comment by ozgrakkurt 3 days ago
Comment by dekhn 3 days ago
Who cares how big the graph is in CSV? That's not the representation you operate over in big data.
All of this would have easily fit in memory on any reasonable modern system.
Comment by adsharma 3 days ago
Comment by ssinchenko 2 days ago
Understand me correctly. This is my research project and I only have a laptop, not a server with 256 GB of RAM. I tested my project on a 2B graph with a hard cap of 8–10 GB of RAM. Of course it fits in memory on any modern system with 64–128 GB of RAM. The whole idea was to conduct a stress test and check how my tool works in out-of-core mode, not to prove to anyone that 2 billion edges (30 GB CSV) constitutes "big data".
Comment by slopblast 3 days ago
Comment by convolvatron 3 days ago
Comment by esafak 3 days ago
Comment by ssinchenko 2 days ago
The multi-processing relies on DataFusion's Tokyo workers. The out-of-core aspect is achieved through a combination of DataFusion FairSpillPool, Sort-Merge-Join, and manually offloading everything to temporary Parquet files on disk.
Comment by Natalia724 3 days ago
Comment by lmeyerov 3 days ago
The cool in the original post was directly inspired by our work here, with our advocacy to the author of keeping their previous Spark work for initial data lake data extraction, and the actual graph work to be redone in our columnar in-memory optimized style for magnitudes of speedup , cost savings
Pip install, benchmarks : https://pygraphistry.readthedocs.io/en/latest/gfql/benchmark...