Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
Next Next commit
[WIP] prepare the release
  • Loading branch information
SemyonSinchenko committed Oct 7, 2025
commit 6da6044295d5697240c30e6e46814b5e78d192cf
8 changes: 5 additions & 3 deletions docs/src/02-quick-start/02-quick-start.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,13 +86,15 @@ a conversion method (see [user guide](/04-user-guide/12-graphx-conversion.md)).
supported algorithms:

| Algorithm | GraphX Wrapper | GraphFrames Implementation | Recommendations |
|--------------------------------|----------------|----------------------------|--------------------------------------------------------------|
| ------------------------------ | -------------- | -------------------------- | ------------------------------------------------------------ |
| BFS | Yes | Yes | GraphFrames provides smoother API |
| Connected Components | Yes | Yes | For small graphs and streaming GraphX, otherwise GraphFrames |
| Strongly Connected Components | Yes | No | GraphX |
| Label Propagation Algorithm | Yes | Yes | GraphFrames is order of magnitude faster |
| Label Propagation Algorithm | Yes | Yes | For small graphs and streaming GraphX, otherwise GraphFrames |
| PageRank | Yes | No | GraphX |
| Parallel Personalized PageRank | Yes | No | GraphX |
| Shortest Paths | Yes | Yes | For small graphs and streaming GraphX, otherwise GraphFrames |
| Triangle Count | Yes | Yes | GraphFrames provides smoother API |
| SVD++ | Yes | No | GraphX |
| SVD++ | Yes | No | GraphX |
| Cycles Detection | No | Yes | GraphFrames |
| Triangel Count | No | Yes | GraphFrames |
5 changes: 5 additions & 0 deletions docs/src/03-tutorials/01-tutorials.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,3 +13,8 @@ More tutorials will be added in future releases to cover additional GraphFrames
*These tutorials are not part of the official GraphFrames documentation.*

- [GraphFrames on Databricks](https://docs.databricks.com/aws/en/integrations/graphframes/)
- [Advanced Deduplication Using Apache Spark: A Guide for Machine Learning Pipelines](https://dev.to/ajay_tani/advanced-deduplication-using-apache-spark-a-guide-for-machine-learning-pipelines-3ik6)
- [Parallelized Community Detection in Social Networks](https://github.com/kedarghule/Community-Detection-in-Social-Networks)
- [Determining Communities In Graphs Using Label Propagation Algorithm](https://hackernoon.com/determining-communities-in-graphs-using-label-propagation-algorithm-4f3m31mm)
- [Top 10 Graph Processing Algorithms in Spark - A Comprehensive Guide for Data Scientists](https://moldstud.com/articles/p-top-10-graph-processing-algorithms-in-spark-a-comprehensive-guide-for-data-scientists)
- [Exploring GraphFrames in PySpark](https://medium.com/tomtalkspython/exploring-graphframes-in-pyspark-f948c9a39844)
161 changes: 153 additions & 8 deletions docs/src/04-user-guide/05-traversals.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,14 @@ vertex ID. Note that this takes an edge direction into account.

See [Wikipedia](https://en.wikipedia.org/wiki/Shortest_path_problem) for a background.

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

### Python API

For API details, refer to the @:pydoc(graphframes.GraphFrame.shortestPaths).
Expand All @@ -33,6 +41,28 @@ val results = g.shortestPaths.landmarks(Seq("a", "d")).run()
results.select("id", "distances").show()
```

### Arguments

- `landmarks`

The list (`Seq`) of vertices that are used as landmarks to compute shortest paths from them to all other vertices.

- `algorithm`

Possible values are `graphx` and `graphframes`. Both implementations are based on the same logic. GraphX is faster for small-medium sized graphs but requires more memory due to less efficient RDD serialization and it's triplets-based nature. GraphFrames requires much less memory due to efficient Thungsten serialization and because the core structures are edges and messages, not triplets.

- `checkpoint_interval`

For `graphframes` only. To avoid exponential growing of the Spark' Logical Plan, DataFrame lineage and query optimization time, it is required to do checkpointing periodically. While checkpoint itself is not free, it is still recommended to set this value to something less than `5`.

- `use_local_checkpoints`

For `graphframes` only. By default, GraphFrames uses persistent checkpoints. They are realiable and reduce the errors rate. The downside of the persistent checkpoints is that they are requiride to set up a `checkpointDir` in persistent storage like `S3` or `HDFS`. By providing `use_local_checkpoints=True`, user can say GraphFrames to use local disks of Spark' executurs for checkpointing. Local checkpoints are faster, but they are less reliable: if the executur lost, for example, is taking by the higher priority job, checkpoints will be lost and the whole job fails.

- `storage_level`

The level of storage for intermediate results and the output `DataFrame` with components. By default it is memory and disk deserialized as a good balance between performance and reliability. For very big graphs and out-of-core scenarious, using `DISK_ONLY` may be faster.

## Breadth-first search (BFS)

Breadth-first search (BFS) finds the shortest path(s) from one vertex (or a set of vertices) to another vertex (or a set
Expand Down Expand Up @@ -88,12 +118,17 @@ Computes the connected component membership of each vertex and returns a graph w

See [Wikipedia](https://en.wikipedia.org/wiki/Connected_component_(graph_theory)) for the background.

**NOTE:** With GraphFrames 0.3.0 and later releases, the default Connected Components algorithm requires setting a Spark
checkpoint directory. Users can revert to the old algorithm using `connectedComponents.setAlgorithm("graphx")`. Starting
from GraphFrames 0.9.3 release, users can also use `localCheckpoints` that does not require setting a Spark checkpoint
directory. To use `localCheckpoints` users can set the config `spark.graphframes.useLocalCheckpoints` to `true` or use
the API `connectedComponents.setUseLocalCheckpoints(true)`. While `localCheckpoints` provides better performance they
are not as reliable as the persistent checkpointing.
---

**NOTE:**

*With GraphFrames 0.3.0 and later releases, the default Connected Components algorithm requires setting a Spark checkpoint directory. Users can revert to the old algorithm using `connectedComponents.setAlgorithm("graphx")`. Starting from GraphFrames 0.9.3 release, users can also use `localCheckpoints` that does not require setting a Spark checkpoint directory. To use `localCheckpoints` users can set the config `spark.graphframes.useLocalCheckpoints` to `true` or use the API `connectedComponents.setUseLocalCheckpoints(true)`. While `localCheckpoints` provides better performance they are not as reliable as the persistent checkpointing.*

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

### Python API

Expand Down Expand Up @@ -123,13 +158,59 @@ val result = g.connectedComponents.setUseLocalCheckpoints(true).run()
result.select("id", "component").orderBy("component").show()
```

### Arguments

- `algorithm`

Possible values are `graphx` and `graphframes`. GraphX-based implementation is a pure Pregel one-by-one. While it may be slightly faster on small-medium sized graphs, it has a much bigger convergence complexity and requires much more memory due to less efficient RDD serialization. GraphFrame-based implementation is based on the ideas from the [Kiveris, Raimondas, et al. "Connected components in mapreduce and beyond." Proceedings of the ACM Symposium on Cloud Computing. 2014.](https://dl.acm.org/doi/abs/10.1145/2670979.2670997). This implementation has much better convergence complexity as well as requires less amount of memory.

- `maxIter`

For `graphx` only. Limit the maximal amount of Pregel iterations. By default it is infinity (`Integer.maxValue`). It is recommended do not change this value. If the algorithm stucks, it is a problem of the graph, not algorithm.

- `checkpoint_interval`

For `graphframes` only. To avoid exponential growing of the Spark' Logical Plan, DataFrame lineage and query optimization time, it is required to do checkpointing periodically. While checkpoint itself is not free, it is still recommended to set this value to something less than `5`.

- `broadcast_threshold`

For `graphframes` only. See [this section](05-traversals.md#aqe-broadcast-mode) for details.

- `use_labels_as_components`

For `graphframes` only. In the case, when the type of the input graph vertices is not one of `Long`, `Int`, `Short`, `Byte`, output labels (components) are a random `Long` numbers by default. By providing `use_labels_as_components=True` user can ask GraphFrames to use original vertex labels for output components. In that case, the minimal value of all original IDs will be used for each of found components. This operation is not free and require an additional `groupBy` + `agg` + `join`.

- `use_local_checkpoints`

For `graphframes` only. By default, GraphFrames uses persistent checkpoints. They are realiable and reduce the errors rate. The downside of the persistent checkpoints is that they are requiride to set up a `checkpointDir` in persistent storage like `S3` or `HDFS`. By providing `use_local_checkpoints=True`, user can say GraphFrames to use local disks of Spark' executurs for checkpointing. Local checkpoints are faster, but they are less reliable: if the executur lost, for example, is taking by the higher priority job, checkpoints will be lost and the whole job fails.

- `storage_level`

The level of storage for intermediate results and the output `DataFrame` with components. By default it is memory and disk deserialized as a good balance between performance and reliability. For very big graphs and out-of-core scenarious, using `DISK_ONLY` may be faster.

### AQE-broadcast mode

*Starting from 0.10.0*

For `graphframes` algorithm only. During iterations, this algorithm can generate new edges that may tend to high skewness in joins and aggregates, because some vertices are having a very-high degree. In previous versions of GraphFrames this issue was addressed by manual broadcasting very high-degree nodes. Unfortunately, Apache Spark Adaptive Quey Execution optimization fails on such a case and that was the reason shy AQE was disabled for Connected Components.

In the new versions of GraphFrames (0.10+) there is a way to disable manual broadcasting, enable AQE and allow it to handle skewnewss. To enable this mode, pass `-1` to the `setBroadcastThreshold`. Based on benchmarks, this mode provides about 5x speed-up. It is possible, that in the future releases, the default value of the `broadcastThreshold` will be changed to `-1`.

### Strongly connected components

Compute the strongly connected component (SCC) of each vertex and return a graph with each vertex assigned to the SCC
containing that vertex. At the moment, SCC in GraphFrames is a wrapper around GraphX implementation.

See [Wikipedia](https://en.wikipedia.org/wiki/Strongly_connected_component) for the background.

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

#### Python API

For API details, refer to the @:pydoc(graphframes.GraphFrame.stronglyConnectedComponents).
Expand Down Expand Up @@ -162,6 +243,14 @@ result.select("id", "component").orderBy("component").show()

Computes the number of triangles passing through each vertex.

---
**WARNING!**


*The current implementation is based on collecting neighbor sets for vertices and compute a pairwise intersection of them. While this works for regular graphs, it will most probably fail on any kind of power-law graphs (graphs with a few very high-degree vertices) or at least will require a lot of memory for Spark Cluster. Consider edge sampling strategies before running the algorithm to get an approximate count of triangles.*

---

### Python API

For API details, refer to the @:pydoc(graphframes.GraphFrame.triangleCount).
Expand Down Expand Up @@ -193,6 +282,51 @@ results.select("id", "count").show()
GraphFrames provides an implementation of
the [Rocha–Thatte cycle detection algorithm](https://en.wikipedia.org/wiki/Rocha%E2%80%93Thatte_cycle_detection_algorithm).

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

**WARNING:**

- *This algorithm returns all the cycles, and users should handle deduplication of \[1, 2, 1\] and \[2, 1, 2\] (that is the same cycle)!*
- *This algorithm collects the full sequences and may require a lot of cluste memory for power-law graphs*

---

### Python API

```python
from graphframes import GraphFrame

vertices = spark.createDataFrame(
[(1, "a"), (2, "b"), (3, "c"), (4, "d"), (5, "e")], ["id", "attr"]
)
edges = spark.createDataFrame(
[(1, 2), (2, 3), (3, 1), (1, 4), (2, 5)], ["src", "dst"]
)
graph = GraphFrame(vertices, edges)
res = graph.detectingCycles(
checkpoint_interval=3,
use_local_checkpoints=True,
)
res.show(False)

# Output:
# +----+--------------+
# | id | found_cycles |
# +----+--------------+
# |1 |[1, 3, 1] |
# |1 |[1, 2, 1] |
# |1 |[1, 2, 5, 1] |
# |2 |[2, 1, 2] |
# |2 |[2, 5, 1, 2] |
# |3 |[3, 1, 3] |
# |5 |[5, 1, 2, 5] |
# +----+--------------+
```

### Scala API

```scala
Expand Down Expand Up @@ -222,5 +356,16 @@ res.show(false)
// +----+--------------+
```

**WARNING:** This algorithm returns all the cycles, and users should handle deduplication of \[1, 2, 1\] and \[2, 1, 2\] (
that is the same cycle)!
### Arguments

- `checkpoint_interval`

For `graphframes` only. To avoid exponential growing of the Spark' Logical Plan, DataFrame lineage and query optimization time, it is required to do checkpointing periodically. While checkpoint itself is not free, it is still recommended to set this value to something less than `5`.

- `use_local_checkpoints`

For `graphframes` only. By default, GraphFrames uses persistent checkpoints. They are realiable and reduce the errors rate. The downside of the persistent checkpoints is that they are requiride to set up a `checkpointDir` in persistent storage like `S3` or `HDFS`. By providing `use_local_checkpoints=True`, user can say GraphFrames to use local disks of Spark' executurs for checkpointing. Local checkpoints are faster, but they are less reliable: if the executur lost, for example, is taking by the higher priority job, checkpoints will be lost and the whole job fails.

- `storage_level`

The level of storage for intermediate results and the output `DataFrame` with components. By default it is memory and disk deserialized as a good balance between performance and reliability. For very big graphs and out-of-core scenarious, using `DISK_ONLY` may be faster.
38 changes: 38 additions & 0 deletions docs/src/04-user-guide/06-graph-clustering.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,14 @@ Run a static Label Propagation Algorithm for detecting communities in networks.

See [Wikipedia](https://en.wikipedia.org/wiki/Label_Propagation_Algorithm) for the background.

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

### Python API

For API details, refer to the @:pydoc(graphframes.GraphFrame.labelPropagation).
Expand All @@ -32,10 +40,40 @@ val result = g.labelPropagation.maxIter(5).run()
result.select("id", "label").show()
```

### Arguments

- `maxIter`

An amount of Pregel iterations. While in theory, Label Propagation algorithm should converge sooner or later to some stable state, there are a lot of problems with it on a real-world graphs. The first one is oscillations: even if the algorithm is almost converged, on a big graphs some vertices at the border between detected communities may contibue oscilate from one iteration to another. The biggest problme, however, is that algorithm may easily converge to the state when all vertices has the same label. It is strongly recommended to set `maxIter` to some reasonable value from `5` to `10` and do some experiments depends of the task and the goal.

- `algorithm`

Possible values are `graphx` and `graphframes`. Both implementations are based on the same logic. GraphX is faster for small-medium sized graphs but requires more memory due to less efficient RDD serialization and it's triplets-based nature. GraphFrames requires much less memory due to efficient Thungsten serialization and because the core structures are edges and messages, not triplets.

- `checkpoint_interval`

For `graphframes` only. To avoid exponential growing of the Spark' Logical Plan, DataFrame lineage and query optimization time, it is required to do checkpointing periodically. While checkpoint itself is not free, it is still recommended to set this value to something less than `5`.

- `use_local_checkpoints`

For `graphframes` only. By default, GraphFrames uses persistent checkpoints. They are realiable and reduce the errors rate. The downside of the persistent checkpoints is that they are requiride to set up a `checkpointDir` in persistent storage like `S3` or `HDFS`. By providing `use_local_checkpoints=True`, user can say GraphFrames to use local disks of Spark' executurs for checkpointing. Local checkpoints are faster, but they are less reliable: if the executur lost, for example, is taking by the higher priority job, checkpoints will be lost and the whole job fails.

- `storage_level`

The level of storage for intermediate results and the output `DataFrame` with components. By default it is memory and disk deserialized as a good balance between performance and reliability. For very big graphs and out-of-core scenarious, using `DISK_ONLY` may be faster.

## Power Iteration Clustering (PIC)

GraphFrames provides a wrapper for the [Power Iteration Clustering](https://www.cs.cmu.edu/~frank/papers/icml2010-pic-final.pdf) algorithm from the SparkML library.

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

### Python API

```python
Expand Down
10 changes: 9 additions & 1 deletion docs/src/04-user-guide/09-aggregate-messages.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,14 @@ Like GraphX, GraphFrames provides primitives for developing graph algorithms. Th
* `aggregateMessages`: Send messages between vertices, and aggregate messages for each vertex. GraphFrames provides a native `aggregateMessages` method implemented using DataFrame operations. This may be used analogously to the GraphX API.
* joins: Join message aggregates with the original graph. GraphFrames rely on `DataFrame` joins, which provide the full functionality of GraphX joins.

---

**NOTE**

*Be aware, that returned `DataFrame` is persistent and should be unpersisted manually after processing to avoid memory leaks!*

---

The below code snippets show how to use `aggregateMessages` to compute the sum of the ages of adjacent users.

## Python API
Expand Down Expand Up @@ -52,4 +60,4 @@ val agg = { g.aggregateMessages
agg.show()
```

For a more complex example, look at the code used to implement the @:pydoc(graphframes.examples.BeliefPropagation).
For a more complex example, look at the code used to implement the @:pydoc(graphframes.examples.BeliefPropagation).
Loading