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.

Advertisement

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

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

Advertisement

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 else

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

A motif is a join plan: each edge term joins the edge table to what came beforePattern(a)-[e1]->(b); (b)-[e2]->(c); (c)-[e3]->(a)Vertices DataFrame (id, ...)Edges DataFrame (src, dst, ...)edges AS e1 (a to b)edges AS e2 (b to c)edges AS e3 (c to a)JOIN e1.dst = e2.srcrows: all 2-hop pathsJOIN e3.src = e2.dst AND e3.dst = e1.srcrows: directed trianglesJOIN vertices for a, b, cattach vertex columnsIntermediate size grows with every hop: 2-hop paths can far outnumber edges when hubs exist.Filter early, drop hubs, and check the plan with explain() before running on the full graph.
A three-edge motif runs as successive self-joins of the edge DataFrame; intermediate path counts grow quickly around hub vertices, so filter before matching.

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 column

bfs 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

CallReturnsKnobs that matter
connectedComponents()component per vertexalgorithm ('graphframes' default or 'graphx'), checkpointInterval (default 2), broadcastThreshold (default 1,000,000); needs a checkpoint directory
stronglyConnectedComponents(maxIter)component per vertexmaxIter
pageRank(resetProbability=0.15, ...)GraphFrame with pagerankmaxIter or tol; sourceId for personalised PageRank
labelPropagation(maxIter)label per vertexmaxIter; results vary between runs on ties
triangleCount(storage_level)count per vertexPython requires a StorageLevel; algorithm defaults to 'exact'
shortestPaths(landmarks)distances mapKeep 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 maxIter or 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 setCheckpointDir before 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 rotations

The 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

  1. Install the io.graphframes artifact that matches your Spark and Scala versions, and graphframes-py if you use Python.
  2. Build vertices from edges, deduplicate IDs and look at the degree distribution before any analysis.
  3. Decide how to treat hubs, cap, drop or keep, and write the rule down.
  4. Set a checkpoint directory and cache the graph before running iterative algorithms.
  5. Start motif queries on a filtered subgraph and inspect explain() before scaling up.
  6. Put explicit iteration limits on PageRank, label propagation and Pregel jobs, and log the counts.
  7. Validate results on a small hand-checked graph, then compare component-size distributions run to run.
Key takeaway: GraphFrames turns a graph into two DataFrames, so graph work uses Spark SQL's optimizer, formats and pipelines. Motifs are self-joins, traversals and Pregel are repeated joins and aggregations, and the built-in algorithms iterate the same way. Install the io.graphframes artifact that matches your cluster, derive vertices from edges, deal with hubs before they deal with you, set a checkpoint directory, cache the graph, bound iterations and filter before pattern matching.