Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 14 additions & 14 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -98,13 +98,13 @@ g.degrees.show()
g2 = g.pageRank(resetProbability=0.15, tol=0.01)
g2.vertices.show()

# +---+-----+---+------------------+
# | id| name|age| pagerank|
# +---+-----+---+------------------+
# | 1| John| 30|0.7758750474847483|
# | 2|Alice| 25|1.4482499050305027|
# | 3| Bob| 35|0.7758750474847483|
# +---+-----+---+------------------+
# +---+-------+---+------------------+
# | id| name|age| pagerank|
# +---+-------+---+------------------+
# | 1| Alice| 30|0.7758750474847483|
# | 2| Bob| 25|1.4482499050305027|
# | 3|Charlie| 35|0.7758750474847483|
# +---+-------+---+------------------+

# GraphFrames' most used feature...
# Connected components can do big data entity resolution on billions or even trillions of records!
Expand All @@ -113,13 +113,13 @@ g2.vertices.show()
sc.setCheckpointDir("/tmp/graphframes-example-connected-components") # required by GraphFrames.connectedComponents
g.connectedComponents().show()

# +---+-----+---+---------+
# | id| name|age|component|
# +---+-----+---+---------+
# | 1| John| 30| 1|
# | 2|Alice| 25| 1|
# | 3| Bob| 35| 1|
# +---+-----+---+---------+
# +---+-------+---+---------+
# | id| name|age|component|
# +---+-------+---+---------+
# | 1| Alice| 30| 1|
# | 2| Bob| 25| 1|
# | 3|Charlie| 35| 1|
# +---+-------+---+---------+

# Find frenemies with network motif finding! See how graph and relational queries are combined?
(
Expand Down
35 changes: 20 additions & 15 deletions docs/mdoc/01-installation.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ For Spark 4.x:
./sbin/start-connect-server.sh \
--conf spark.connect.extensions.relation.classes=\
org.apache.spark.sql.graphframes.GraphFramesConnect \
--packages io.graphframes.graphframes-connect-spark4_2.13:@VERSION@
--packages io.graphframes:graphframes-connect-spark4_2.13:@VERSION@
```

For Spark 3.x:
Expand All @@ -90,7 +90,7 @@ For Spark 3.x:
./sbin/start-connect-server.sh \
--conf spark.connect.extensions.relation.classes=\
org.apache.spark.sql.graphframes.GraphFramesConnect \
--packages io.graphframes.graphframes-connect-spark3_2.12:@VERSION@
--packages io.graphframes:graphframes-connect-spark3_2.12:@VERSION@
```

**WARNING**: The GraphFrames Connect Server Extension is not compatible with managed SparkConnect from Databricks. To make it work, you need to use build GraphFrames Connect Server Extension from source with a flag:
Expand All @@ -116,19 +116,24 @@ message GraphFramesAPI {
BFS bfs = 4;
ConnectedComponents connected_components = 5;
DropIsolatedVertices drop_isolated_vertices = 6;
FilterEdges filter_edges = 7;
FilterVertices filter_vertices = 8;
Find find = 9;
LabelPropagation label_propagation = 10;
PageRank page_rank = 11;
ParallelPersonalizedPageRank parallel_personalized_page_rank = 12;
PowerIterationClustering power_iteration_clustering = 13;
Pregel pregel = 14;
ShortestPaths shortest_paths = 15;
StronglyConnectedComponents strongly_connected_components = 16;
SVDPlusPlus svd_plus_plus = 17;
TriangleCount triangle_count = 18;
Triplets triplets = 19;
DetectingCycles detecting_cycles = 7;
FilterEdges filter_edges = 8;
FilterVertices filter_vertices = 9;
Find find = 10;
LabelPropagation label_propagation = 11;
PageRank page_rank = 12;
ParallelPersonalizedPageRank parallel_personalized_page_rank = 13;
PowerIterationClustering power_iteration_clustering = 14;
Pregel pregel = 15;
ShortestPaths shortest_paths = 16;
StronglyConnectedComponents strongly_connected_components = 17;
SVDPlusPlus svd_plus_plus = 18;
TriangleCount triangle_count = 19;
Triplets triplets = 20;
KCore kcore = 21;
MaximalIndependentSet mis = 22;
RandomWalkEmbeddings rw_embeddings = 23;
AggregateNeighbors aggregate_neighbors = 24;
}
}
```
Expand Down
2 changes: 1 addition & 1 deletion docs/src/01-about/01-index.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ GraphFrames should be compatible with any platform that runs the open-source Spa

GraphFrames is compatible with Spark 3.4+. However, later versions of Spark include major improvements to DataFrames, so GraphFrames may be more efficient when running on more recent Spark versions.

GraphFrames is tested with Java 8, 11 and 17, Python 3, Spark 3.5 and Spark 4.0 (Scala 2.12 / Scala 2.13).
GraphFrames is tested with Java 8, 11 and 17, Python 3.10-3.13, Spark 3.5, Spark 4.0, and Spark 4.1 (Scala 2.12 / Scala 2.13).

# Applications, the Apache Spark shell, and clusters

Expand Down
4 changes: 2 additions & 2 deletions docs/src/01-about/02-architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,8 @@ Let’s look at a concrete example – PageRank. This algorithm became famous fo
title = "PageRank Algorithm"
}

In GraphFrames, most algorithms – including PageRank – are built on the Pregel framework ([*Malewicz, Grzegorz, et al. "Pregel: a system for large-scale graph processing." Proceedings of the 2010 ACM SIGMOD International Conference on Management of data. 2010.*](https://blog.lavaplanets.com/wp-content/uploads/2023/12/p135-malewicz.pdf)). We represent the graph as two `DataFrames`, which you can think of as tables: one for edges and one for vertices. The PageRank table is initialized by assigning every vertex a starting rank of `0.15`.
In GraphFrames, many algorithms are built on the Pregel framework ([*Malewicz, Grzegorz, et al. "Pregel: a system for large-scale graph processing." Proceedings of the 2010 ACM SIGMOD International Conference on Management of data. 2010.*](https://blog.lavaplanets.com/wp-content/uploads/2023/12/p135-malewicz.pdf)). Some algorithms, such as PageRank, currently rely on GraphX-backed implementations. We represent the graph as two `DataFrames`, which you can think of as tables: one for edges and one for vertices. The PageRank table is initialized by assigning every vertex a starting rank of `1.0`.

Each iteration of PageRank works like a series of SQL operations. The process starts by joining the edges table with the current PageRank values for each vertex. This creates a triplets table, where each row contains a source, destination, and their current ranks. Next, we generate messages: each source sends its rank to its destination. These messages are grouped by destination and summed up. Finally, we join the results back to the PageRank table and update the rank using a simple formula: `new_rank = sum_rank * 0.85 + 0.15`
Each iteration of PageRank works like a series of SQL operations. The process starts by joining the edges table with the current PageRank values for each vertex. This creates a triplets table, where each row contains a source, destination, and their current ranks. Next, we generate messages: each source sends its rank to its destination. These messages are grouped by destination and summed up. Finally, we join the results back to the PageRank table and update the rank using a simple formula: `new_rank = sum_rank * 0.85 + 0.15`, where `0.85` is the damping factor and `0.15` is the reset probability (alpha) that vertices with no in-links converge to.

This whole process is repeated – each step is just a combination of joins, group by, and aggregates over tables – until the ranks stop changing much. The algorithm converges quickly, usually in about 15–20 iterations. Since it relies entirely on SQL operations, running PageRank on an Apache Spark cluster gives you excellent horizontal scalability. As long as your tables fit in Spark, you can compute PageRank using Pregel. In practice, this means you can almost infinitely scale just by adding more hardware.
6 changes: 3 additions & 3 deletions docs/src/01-about/03-benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,8 @@
## Graphalytics Benchmarks

This benchmark is to test the performance of GraphFrames algorithms, not Apache Spark itself. So, all the graphs are
read from the disk and persisted in memory in the serialized format. In the result, only the time of GraphFrames
algorithms is measured and the time of reading of the CSV, serialization and persisting the data does not measure.
read from Parquet files on disk and persisted in memory in the serialized format. As a result, only the time of GraphFrames
algorithms is measured, and the time to read/parse source files, serialize, and persist the data is not measured.

### Configurations

Expand All @@ -19,7 +19,7 @@ algorithms is measured and the time of reading of the CSV, serialization and per
- **Vertices:** 2M
- **Edges:** 5M
- **Size Category:** _XS_
- **Source files format:** `CSV`-like
- **Source files format:** `Parquet`

| Algorithm | Measurements | Time (s) |
| -------------------------------- | ---------------------------------------------- | ---------------------------------------- |
Expand Down
16 changes: 8 additions & 8 deletions docs/src/06-contributing/01-contributing-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,12 +14,12 @@ Ensure the following tools are installed before cloning the repository:
| --- | --- | --- |
| Git | Latest stable | Required for version control and contribution workflows. |
| Java Development Kit (JDK) | 11 or 17 | Spark 3.x supports Java 8/11/17; GraphFrames CI runs on JDK 17. |
| Python | 3.10 – 3.12 | Required for the Python APIs and tests. |
| Apache Spark (binary distribution) | 3.5.x (default) or 4.0.x | Needed for the Python test suite. |
| Python | 3.10 – 3.13 | Required for the Python APIs and tests. |
| Apache Spark (binary distribution) | 3.5.x (default), 4.0.x, or 4.1.x | Needed for the Python test suite. |
| Poetry | ≥ 1.8 | Dependency manager used by the Python package. Install via [`pipx`](https://pypa.github.io/pipx/) or `pip`. |
| Protocol Buffers compiler (`protoc`) | ≥ 3.21 | Required for the GraphFrames Connect protobuf build. |
| Buf CLI | Latest stable | Used to lint and generate protobuf sources. |
| Apache Spark (optional) | 3.5.x (default) or 4.0.x | Only required if you want the standalone Spark shell outside PySpark. |
| Apache Spark (optional) | 3.5.x (default), 4.0.x, or 4.1.x | Only required if you want the standalone Spark shell outside PySpark. |
Comment thread
SemyonSinchenko marked this conversation as resolved.
| Docker (optional) | Latest stable | Useful for isolated environments but not mandatory. |

### 1.1 Install required tooling
Expand Down Expand Up @@ -55,12 +55,12 @@ export PATH="$JAVA_HOME/bin:$PATH"
#### Optional: Standalone Apache Spark distribution
`poetry install` (described later) already brings in the matching version of PySpark and Spark
Connect. If you also want the standalone Spark shell or `spark-submit`, download the distribution
that matches the build’s `spark.version` (currently 3.5.6) and expose it via `SPARK_HOME`:
that matches the build’s `spark.version` (currently 3.5.7) and expose it via `SPARK_HOME`:
```bash
curl -O https://downloads.apache.org/spark/spark-3.5.6/spark-3.5.6-bin-hadoop3.tgz
curl -O https://downloads.apache.org/spark/spark-3.5.7/spark-3.5.7-bin-hadoop3.tgz
mkdir -p "$HOME/.local/spark"
tar -xzf spark-3.5.6-bin-hadoop3.tgz -C "$HOME/.local/spark"
export SPARK_HOME="$HOME/.local/spark/spark-3.5.6-bin-hadoop3"
tar -xzf spark-3.5.7-bin-hadoop3.tgz -C "$HOME/.local/spark"
export SPARK_HOME="$HOME/.local/spark/spark-3.5.7-bin-hadoop3"
export PATH="$SPARK_HOME/bin:$PATH"
```

Expand Down Expand Up @@ -302,7 +302,7 @@ To build documentation and run a preview server run `./build/sbt docs/laikaPrevi
| Run a specific Scala suite | `./build/sbt "core/testOnly <SuiteName>"` |
| Build assembly jar | `./build/sbt core/assembly` |
| Install Python dependencies | `cd python && poetry install --with dev` |
| Run Python tests | `cd python && ./run-tests.sh` |
| Run Python tests | `cd python && poetry run pytest -vvv` |
| Run Python formatters | `poetry run black graphframes tests` |
| Install pre-commit hooks | `pre-commit install` |

Expand Down
Loading