Algorithms on billion-scale graph using 10GB RAM: I love DataFusion

(semyonsinchenko.github.io)

71 points | by speckx 4 hours ago

13 comments

  • chrisweekly 3 hours ago
    > "I can compute PageRank on a directed graph with one billion edges (graph500-26 from the Graphalytics dataset) using 5 GB of memory. Alternatively, I can identify all the weakly connected components in a graph with two billion edges (twitter_mpi from the same dataset collection) using 10 GB of memory. Neither NetworkX nor Igraph can do this; most existing graph algorithms require the graph to fit into memory. Previously, I thought you needed Apache Spark and GraphFrames for billion-scale graph analytics. Now, however, I think all you need is a laptop. I have completely changed my old opinion about using Apache DataFusion for graph analytics."

    Impressive!

    • danbruc 2 hours ago
      33M vertices and 1B edges easily fits into memory, if you use 32 bit integers, it will require about 4 GiB of memory. Out of curiosity I just implemented generating a random graph of that size and calculating one page rank iteration on it in the most naive way (20 lines of C#) and it consumed 8.4 GiB of memory and got one iteration done in 4:20 minutes single threaded.
      • anon7725 3 hours ago
        > most existing graph algorithms require the graph to fit into memory.

        You can get pretty far with sparse graphs, which are just arrays, in combination with memory mapping.

      • yadgire7 2 hours ago
        Hello, I am new to hacker news and finding it really resourceful. I found this article interesting (having learnt KG and Map Reduce (spark) as part of my masters' course), appreciate the effort to post this.

        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!

        • RobinL 29 minutes ago
          Love this Sem. Will provide a great alternative to the SQL based connected components algorithm we ship in Splink. Looking forward to testing how much faster it is. Thanks for your work on it!
          • cpdomina 1 hour ago
            cool! you might be interested in graphchi (2012), also designed to do large scale graph operations on a single machine

            https://github.com/GraphChi/graphchi-cpp#performance

            • adsharma 1 hour ago
              It looks similar to networkit in that it represents graphs using row oriented memory layout.

              I forked networkit for exactly this reason. Columnar memory is much more efficient.

            • lmeyerov 1 hour ago
              Related, we recently release the polars version of GFQL, the only oss cypher property graph query engine for CPU+GPU, and even better, no database nor outside process needed. We started doing LDBC benchmarks vs neo4j, memgraph, kuzu, etc, and are already starting to outperform them both on latency for small OLTP graph searches and $, speed for big graph OLAP ones, especially in GPU mode.

              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...

              • adsharma 1 hour ago
                The idea of graph algorithms on Apache arrow at scale originated here. 100+ graph algorithms running on columnar memory.

                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.

                • lmeyerov 4 minutes ago
                  Not really ;-)

                  ~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.

                  • adsharma 1 hour ago
                    https://github.com/LadybugDB/ladybug-icebug-notebooks/blob/m...

                    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.

                  • ratmice 3 hours ago
                    It would be nice if OP noted what caused the change in their opinion?

                    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?

                    • lmeyerov 1 minute ago
                      We shared with the author how databricks graphframes were wildly inefficient for this kind of thing compared to our benchmarking in pygraphistry & gfql: we were measuring doing billion-edge graph traversals & scans in single node in-memory in seconds, so the core of pagerank, which is magnitudes more efficient than their original spark approach.

                      So then the question became pandas/polars/datafusion/duckdb/etc, must of which are rust/native. I'm curious myself why datafusion vs others :)

                    • dekhn 1 hour ago
                      It's hard to take the article seriously when it has quotes like this: "The hardest part. 2B edges twitter graph is already huge (its edges are 30 GB in CSV !!!)."

                      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.

                    • ozgrakkurt 2 hours ago
                      3.4 gb dataset on 10gb ram
                      • theLiminator 3 hours ago
                        DataFusion is really cool, it's kind of like the LLVM of the OLAP world.
                        • slopblast 4 hours ago
                          Really cool visualization, amazing how it resembles a neural network.
                          • convolvatron 2 hours ago
                            I'm pretty sure that's some stock output from CAIDA, looks like a traceroute graph from the inset
                          • esafak 3 hours ago
                            Does it support out-of-core or multi-processor processing?
                            • Natalia724 3 hours ago
                              [dead]