GraphFrames is a graph library for Apache Spark that represents a graph as two DataFrames: one of vertices, one of edges. Because both are ordinary DataFrames, everything Spark SQL can do applies to graphs: read them from Parquet or Delta, filter with SQL expressions, join them to other tables, and let the Catalyst optimizer plan the work. On top of that it adds pattern matching (motif finding), traversals, message passing and a set of standard algorithms, in Python and Scala.
This article explains how GraphFrames actually executes, because that decides which queries finish and which fill the disks with shuffle. It covers installation under the project's newer Maven coordinates, building graphs from tables, motifs as joins, BFS and shortest paths, aggregateMessages and Pregel, the built-in algorithms, a worked fraud-ring example and the operational practices that keep large graph jobs stable. Versions and signatures were checked against graphframes.io, PyPI and Maven Central on 2026-10-03.
The data model
A GraphFrame is built from a vertex DataFrame that must have an id column and an edge DataFrame that must have src and dst columns. Any other columns are properties: a vertex might carry name and country, an edge amount and ts. Edges are directed; an undirected relationship is either stored once and treated as undirected by the algorithm, or stored in both directions.
Every operation returns a DataFrame or a new GraphFrame, so results flow straight into the rest of a Spark pipeline. That is the main difference from GraphX, Spark's older RDD-based graph API, which has its own typed graph representation and is Scala-only. That article also compares the two and graph databases; this one focuses on using GraphFrames well.
Installing it: the coordinates moved
GraphFrames lives outside Apache Spark, so you add it as a package. The project now publishes under the io.graphframes group with one artifact per Spark major version and Scala version, for example graphframes-spark3_2.12 and graphframes-spark4_2.13. At the time of writing the latest release on Maven Central for Spark 3 with Scala 2.12 was 0.12.2. Older tutorials use the legacy Spark Packages coordinates such as graphframes:graphframes:0.8.3-spark3.5-s_2.12; those still resolve but are no longer where new releases go.
# Spark 3.x, Scala 2.12: JVM library on the classpath
pyspark --packages io.graphframes:graphframes-spark3_2.12:0.12.2
# Python wrapper (requires Python 3.10+). It does NOT include the JVM core:
# the cluster, or the Spark Connect server, still needs the package above.
pip install graphframes-pyMatch the artifact to your cluster's Spark and Scala versions exactly; a mismatch usually shows up as a NoSuchMethodError at the first algorithm call, not at import. Recent releases also support Spark Connect: the Python package detects a remote session automatically, and the JVM side must be installed on the Connect server.
Building a graph from tables
from pyspark.sql import functions as F
from graphframes import GraphFrame
payments = spark.read.table("finance.payments") # payer, payee, amount, ts
edges = (payments
.select(F.col("payer").alias("src"), F.col("payee").alias("dst"), "amount", "ts")
.where("src <> dst")) # drop self-loops unless they mean something
vertices = (edges.select(F.col("src").alias("id"))
.union(edges.select(F.col("dst").alias("id")))
.distinct()
.join(spark.read.table("finance.accounts"), "id", "left"))
g = GraphFrame(vertices, edges)
g.inDegrees.orderBy(F.desc("inDegree")).show(10) # find hubs before anything elseThree habits prevent most surprises. Deriving vertices from edges guarantees no edge points at a missing vertex. Removing duplicate vertex IDs matters, because a duplicated ID duplicates every join result that touches it. And looking at the degree distribution first tells you whether there are hubs, vertices with millions of edges, which will dominate the cost of every later step.
Motif finding: patterns become joins
g.find() takes a pattern in a small Cypher-like language. (a)-[e]->(b) names a vertex a, an edge e and a vertex b; terms are separated by semicolons; [] is an anonymous edge; and a term prefixed with ! is a negation, matching only where that edge does not exist. The result is a DataFrame with one struct column per named element, which you then filter with ordinary expressions.
# Directed triangles, then keep only large, fast cycles
tri = g.find("(a)-[e1]->(b); (b)-[e2]->(c); (c)-[e3]->(a)")
suspicious = tri.where("e1.amount > 10000 AND e3.ts - e1.ts < 3600")
# One-way relationships: a pays b, but b never pays a
one_way = g.find("(a)-[]->(b); !(b)-[]->(a)")Under the hood each edge term is a join of the edge DataFrame against the partial result, and vertex columns are joined in after. A three-edge pattern is three self-joins of the edge table. On a graph with hubs the intermediate two-hop result can be orders of magnitude larger than the edge table, which is how a query that looked small fills executor disks. Two remedies work. Filter the graph before matching, with g.filterEdges(...) or g.filterVertices(...), so the joins start smaller. And call .explain() on the result to check the join order and whether broadcasts are used.
Traversals: BFS and shortest paths
# Shortest paths from accounts flagged 'mule' to accounts at 'exchange', max 4 hops
paths = g.bfs(fromExpr="label = 'mule'", toExpr="label = 'exchange'",
edgeFilter="amount > 500", maxPathLength=4)
# Hop distance from every vertex to a few landmark vertices
dist = g.shortestPaths(landmarks=["acct_1", "acct_2"]) # adds a 'distances' map columnbfs returns one row per shortest path found, with the intermediate vertices and edges as columns; the default maxPathLength is 10, and setting it low is the main safeguard. shortestPaths counts hops, not weights. The library also offers an all-paths search, which its own documentation warns can be extremely slow or run out of memory as the length approaches the graph's diameter.
Message passing: aggregateMessages and Pregel
Most graph algorithms are rounds of "every vertex sends a value along its edges, every vertex aggregates what it received". aggregateMessages is one such round, expressed with columns:
from graphframes.lib import AggregateMessages as AM
# Total amount received and sent per account, in one pass over the edges
flows = g.aggregateMessages(
F.sum(AM.msg).alias("total_flow"),
sendToSrc=AM.edge["amount"],
sendToDst=AM.edge["amount"])For iterative algorithms, the pregel builder repeats the round, keeping a vertex column updated each superstep. This is PageRank written by hand, in the shape the project documents:
from graphframes.lib import Pregel
n = g.vertices.count()
alpha = 0.15
ranks = (g.pregel
.setMaxIter(10)
.withVertexColumn("rank", F.lit(1.0 / n),
F.coalesce(Pregel.msg(), F.lit(0.0)) * F.lit(1.0 - alpha) + F.lit(alpha / n))
.sendMsgToDst(Pregel.src("rank") / Pregel.src("outDegree"))
.aggMsgs(F.sum(Pregel.msg()))
.run())This assumes the vertex DataFrame already has an outDegree column, joined from g.outDegrees. Each superstep is a join and an aggregation, so the query plan grows every iteration; that is why checkpointing, covered below, is not optional on long runs.
The built-in algorithms
| Call | Returns | Knobs that matter |
|---|---|---|
connectedComponents() | component per vertex | algorithm ('graphframes' default or 'graphx'), checkpointInterval (default 2), broadcastThreshold (default 1,000,000); needs a checkpoint directory |
stronglyConnectedComponents(maxIter) | component per vertex | maxIter |
pageRank(resetProbability=0.15, ...) | GraphFrame with pagerank | maxIter or tol; sourceId for personalised PageRank |
labelPropagation(maxIter) | label per vertex | maxIter; results vary between runs on ties |
triangleCount(storage_level) | count per vertex | Python requires a StorageLevel; algorithm defaults to 'exact' |
shortestPaths(landmarks) | distances map | Keep the landmark list short |
The DataFrame-based connected components implementation is the workhorse for entity resolution and fraud work. It iterates joins until labels stop changing, checkpoints every few iterations, and treats vertices whose degree exceeds broadcastThreshold specially so a few hubs do not create a skewed shuffle. PageRank itself is covered in depth in Spark PageRank.
Worked example: finding fraud rings
A payments company wants groups of accounts that share devices or cards, because a ring of accounts controlled by one actor usually does. Accounts are vertices; an edge joins two accounts that used the same device fingerprint. Rather than create all pairs per device, which is quadratic for a shared café tablet, the team links each account to the device's first account, which preserves connectivity with one edge per account.
spark.sparkContext.setCheckpointDir("s3://risk-tmp/graphframes-ckpt")
logins = spark.read.table("risk.logins").select("account_id", "device_id").distinct()
# Exclude devices shared by thousands of accounts (public kiosks, emulators' default IDs)
dev_size = logins.groupBy("device_id").count()
logins = logins.join(dev_size.where("count <= 50"), "device_id").drop("count")
anchor = logins.groupBy("device_id").agg(F.min("account_id").alias("anchor"))
edges = (logins.join(anchor, "device_id")
.where("account_id <> anchor")
.select(F.col("account_id").alias("src"), F.col("anchor").alias("dst")))
vertices = logins.select(F.col("account_id").alias("id")).distinct()
cc = GraphFrame(vertices, edges).connectedComponents()
rings = (cc.groupBy("component").agg(F.count("*").alias("size"))
.where("size BETWEEN 5 AND 500"))The device-size cut-off is the decision that makes this work: without it, one emulator ID shared by a hundred thousand accounts merges unrelated people into a single giant component. The team then joins rings to chargeback history, and for the top rings runs the triangle motif above on the payments graph restricted to ring members, which is cheap because the filter shrank the graph from hundreds of millions of edges to a few thousand.
Operating GraphFrames at scale
- Set a checkpoint directory on reliable storage before iterative algorithms. Checkpoints truncate the query plan and lineage; without them, plans grow every iteration until planning time or stack depth becomes the bottleneck.
- Cache what you reuse. The vertex and edge DataFrames are read once per iteration; persisting them, as described in cache and persist, avoids rereading the source tables each time. Unpersist when done.
- Expect skew from hubs. Joins keyed on vertex ID put all of a hub's edges in one task. Remove or cap hubs when the analysis allows, and use adaptive query execution's skew handling; see data skew in Spark.
- Use compact IDs. Long string IDs make every shuffle bigger; mapping them to longs once, up front, pays back on every iteration.
- Bound iterations. Set
maxIteror a tolerance explicitly, and log how many iterations a run used so you notice when the graph's structure changes.
Failure modes
- Executor disks fill during a motif query: intermediate path explosion around hubs. Filter first and check the plan.
- Connected components fails asking for a checkpoint directory: set one with
setCheckpointDirbefore the call. - One giant component: a junk vertex, such as an empty string or a default device ID, is linking everything. Inspect the top-degree vertices.
- Duplicated results: duplicate vertex IDs or duplicate edges multiply join output.
- Version mismatch errors at runtime: the package artifact does not match the cluster's Spark or Scala version.
- Unstable community labels: label propagation can assign different labels across runs; compare membership, not label values.
Testing on a graph you can check by hand
Graph bugs rarely raise errors; they produce plausible wrong numbers. Before running on production data, run the same pipeline on a graph small enough to answer by eye, and keep it as a unit test. The checkpoint directory must already be set, as in the example above.
v = spark.createDataFrame([("a",), ("b",), ("c",), ("d",), ("e",)], ["id"])
e = spark.createDataFrame([("a", "b"), ("b", "c"), ("c", "a"), ("d", "e")], ["src", "dst"])
small = GraphFrame(v, e)
comps = {r["id"]: r["component"] for r in small.connectedComponents().collect()}
assert comps["a"] == comps["b"] == comps["c"] != comps["d"] == comps["e"]
assert small.find("(x)-[]->(y); (y)-[]->(z); (z)-[]->(x)").count() == 3 # one cycle, 3 rotationsThe second assertion encodes a detail worth knowing: a directed triangle matches once per starting vertex, so motif counts must be divided by the pattern's symmetry before reporting them.
What to do next
- Install the
io.graphframesartifact that matches your Spark and Scala versions, andgraphframes-pyif you use Python. - Build vertices from edges, deduplicate IDs and look at the degree distribution before any analysis.
- Decide how to treat hubs, cap, drop or keep, and write the rule down.
- Set a checkpoint directory and cache the graph before running iterative algorithms.
- Start motif queries on a filtered subgraph and inspect
explain()before scaling up. - Put explicit iteration limits on PageRank, label propagation and Pregel jobs, and log the counts.
- Validate results on a small hand-checked graph, then compare component-size distributions run to run.