Connected components answers one question about a graph: which vertices can reach each other if you ignore edge direction? Every vertex gets a component ID, and two vertices share an ID exactly when a path joins them. That sounds like a textbook exercise, yet it is one of the most common graph jobs run on Spark, because it is the core of entity resolution (which accounts belong to one person), fraud ring detection, deduplication of customer records, network segmentation and lineage grouping.
On one machine the problem is easy: a union-find structure processes a billion edges in near-linear time. On a cluster it is not, because a path can wander across every partition, and nothing short of iteration lets the two ends of a long chain learn that they belong together. This article explains the three distributed algorithms you can choose from in Spark today, how many rounds each needs and why, and the operational problems, hubs, query plan growth and unstable IDs, that decide whether the job finishes. Library behaviour is taken from the current GraphX and GraphFrames documentation; option names have changed between GraphFrames releases, so check yours.
Why connected components needs iteration
Start from first principles. A serial union-find keeps a parent pointer per vertex. For each edge (u, v) it finds the root of each endpoint and links one root under the other; with path compression and union by rank every operation is almost constant time. The catch is that find chases pointers that may live anywhere. In a shuffle-based engine, chasing a pointer means a join, so a chain of length k costs k joins unless the algorithm shortens chains as it goes.
All distributed connected-components algorithms are therefore ways of answering one question: how fast can labels or pointers spread through the graph per round, given that one round is a join plus an aggregation over the whole edge set? The answer ranges from the graph's diameter (slow on long chains) to roughly logarithmic in the number of vertices (fast, but each round does more work). If the graph fits in one executor's memory, none of this matters: collect it and run union-find, as covered in union-find. The rest of this article is for graphs that do not.
Min-label propagation in GraphX
The simplest algorithm is min-label propagation, and it is what GraphX's connectedComponents() runs. Every vertex starts with its own VertexId as its label. In each Pregel superstep, every edge whose endpoints carry different labels sends the smaller label to the endpoint with the larger one, and each vertex keeps the minimum it receives. When no messages are sent, every vertex holds the smallest ID in its component.
import org.apache.spark.graphx._
// edges: RDD[Edge[Int]] with Long vertex ids
val graph = Graph.fromEdges(edges, defaultValue = 0)
val cc: VertexRDD[VertexId] = graph.connectedComponents().vertices
// overload with a cap, useful as a guard rather than a correctness tool:
// graph.connectedComponents(maxIterations = 50)
cc.map { case (_, comp) => (comp, 1L) }
.reduceByKey(_ + _)
.sortBy(-_._2)
.take(20) // the 20 largest componentsWorked trace on a path 5 - 3 - 8 - 1 - 9: after round one the labels are 3, 3, 1, 1, 1; after round two 3, 1, 1, 1, 1; after round three every vertex holds 1. A label moves one hop per round, so the number of rounds is bounded by the largest distance a minimum label must travel, roughly the component's diameter. Social graphs with diameter around ten converge quickly. Record-linkage graphs, where a chain of shared addresses can run for hundreds of hops, can take hundreds of supersteps, each one a shuffle. The GraphFrames documentation describes the GraphX path as a naive Pregel-based implementation, which is fair: it is correct and simple, and the right choice only when you know the diameter is small. GraphX internals such as vertex-cut partitioning are covered in Spark GraphX.
Large-star and small-star
Large-star / small-star, from Kiveris and colleagues' 2014 paper "Connected Components in MapReduce and Beyond", attacks the diameter problem by rewiring edges instead of only passing labels. Think of each component as a forest that the algorithm keeps flattening into stars, where every vertex points straight at the component's minimum.
The large-star step: for each vertex u, find the minimum m over u and its neighbours, then replace every edge from u to a neighbour larger than u with an edge from that neighbour to m. The small-star step: for each vertex, find the minimum over its smaller-or-equal neighbours and connect itself and those smaller neighbours to that minimum. Alternating the two shortens long paths geometrically, and the paper proves convergence in a number of rounds logarithmic-squared in the vertex count, with much better behaviour in practice. When the edge set stops changing, every edge joins a vertex to its component's minimum.
# Pseudocode for one round over an undirected edge set E (each edge stored once, u > v)
def large_star(E):
nbrs = group E by endpoint # adjacency including both directions
out = set()
for u, N in nbrs:
m = min(N | {u})
for v in N:
if v > u:
out.add((v, m)) # larger neighbours now point at the min
return out
def small_star(E):
out = set()
for u, N in group E by larger endpoint: # N = neighbours v with v <= u
m = min(N | {u})
for v in N | {u}:
if v != m:
out.add((v, m))
return out
while True:
E2 = small_star(large_star(E))
if E2 == E: break # GraphFrames compares a checksum per round
E = E2
# now every (v, m) has m = min of v's componentRecent GraphFrames documentation calls this algorithm two_phase, describes it as a DataFrame-native implementation based on the large-star / small-star approach, and lists it as the default. Older releases call the same code graphframes, which the documentation now marks as a deprecated alias. Before iterating, the implementation de-duplicates vertices and assigns each a unique Long ID, because comparisons and joins on Longs are much cheaper than on strings.
Randomized contraction
The third option, randomized_contraction, follows Bögeholz, Brand and Todor's ICDE 2020 paper "In-database connected component analysis". Instead of propagating minimum labels, each round maps every vertex through a random linear function and contracts each vertex into the smallest mapped value among itself and its neighbours, collapsing whole neighbourhoods at once. Edges inside a contracted group disappear; the round repeats until no edges remain, then a reverse pass walks the recorded contractions back to assign every original vertex its final component.
According to the GraphFrames documentation, it converges similarly to two_phase in AQE mode, performs better on the project's benchmarks and needs about half the memory. It has one surprise: it always produces random Long component IDs, whatever the input ID type, unless you enable the option the documentation names use_labels_as_components, which instead uses the minimum original vertex label at the cost of an extra group-by, aggregation and join.
Running it in GraphFrames
In Scala the builder exposes the knobs directly. Set a reliable checkpoint directory first; the DataFrame-native algorithms checkpoint periodically and fail without one unless you opt into local checkpoints.
import org.graphframes.GraphFrame
spark.sparkContext.setCheckpointDir("s3a://my-bucket/checkpoints/cc") // HDFS or object store
val g = GraphFrame(vertices, edges) // vertices: id ; edges: src, dst
val cc = g.connectedComponents
.setAlgorithm("two_phase") // or "randomized_contraction", "graphx"
.setCheckpointInterval(2) // docs recommend 2 or below
.setBroadcastThreshold(-1) // two_phase only: let AQE handle skew
.run() // DataFrame: id, ..., component
cc.groupBy("component").count().orderBy(org.apache.spark.sql.functions.desc("count")).show(20)From Python the call is g.connectedComponents() plus keyword arguments for the same settings; keyword names have shifted between graphframes-py releases, so take them from the docstring of the version you install rather than from examples online. The -1 broadcast threshold switches two_phase to an AQE-driven mode that the documentation reports as roughly five times faster than the explicit broadcast path; it requires spark.sql.adaptive.enabled, the default since Spark 3.2. The intermediate storage level defaults to MEMORY_AND_DISK; the documentation suggests DISK_ONLY for very large graphs. Usage beyond connected components is covered in Spark GraphFrames.
| Algorithm | Rounds | Strength | Watch out for |
|---|---|---|---|
| graphx | about the diameter | simple, RDD-based, Pregel | long chains; lineage grows unless spark.graphx.pregel.checkpointInterval is set |
| two_phase | near logarithmic | default, deterministic min-ID labels | edge blow-up around hubs |
| randomized_contraction | near logarithmic | about half the memory per docs | random component IDs by default |
Worked example: one customer ID across 400 million accounts
A marketplace wants one customer ID across 400 million accounts. Two accounts belong together if they share a verified email, a phone number or a payment-card fingerprint. Model it as a bipartite graph: account vertices, key vertices (prefixed so email:a@x.com cannot collide with phone:...), and one edge per account-key pair. Bipartite modelling avoids the quadratic blow-up of writing an edge between every pair of accounts that share a key.
from pyspark.sql import functions as F
links = spark.table("identity.account_keys") # account_id, key_type, key_value
edges = (links
.where(F.col("key_value").isNotNull() & (F.trim("key_value") != ""))
.select(F.col("account_id").alias("src"),
F.concat_ws(":", "key_type", F.lower(F.trim("key_value"))).alias("dst"))
.distinct())
# Hub filter: a key shared by thousands of accounts is a default value, a test card or a reseller, not a person
key_deg = edges.groupBy("dst").count()
hubs = key_deg.where("count > 200").select("dst")
edges = edges.join(hubs, "dst", "left_anti")
vertices = edges.select(F.col("src").alias("id")).union(edges.select(F.col("dst").alias("id"))).distinct()Run connected components on the result and keep only account vertices. The first run showed one component of 31 million accounts. The top-degree query found the cause in minutes: a key value of phone:0000000000 from a signup form default. Adding it to the hub list split the giant into ordinary households and small businesses. Always inspect the size distribution after the first run: a healthy identity graph has a long tail of size-one and size-two components and no single component holding a measurable fraction of all vertices.
The second lesson came the next day. Component IDs are not stable: a new edge can change a component's minimum, and randomized contraction draws fresh random IDs on every run. Downstream tables keyed on yesterday's IDs broke. The fix is a mapping step: join today's components with yesterday's by member, and for each new component inherit the old ID that contributes the most members, minting a new ID only when no old component overlaps. When two old components merge, the larger keeps its ID and the smaller is recorded as an alias, so downstream systems can follow the merge instead of losing history.
Operating the job
Three operational facts dominate run time. First, plan growth: every round of a DataFrame-native algorithm adds joins to the logical plan, and without truncation the optimizer's time grows faster than the work itself. That is why checkpointing exists here; the documentation recommends an interval of two or less. Local checkpoints, enabled with setUseLocalCheckpoints(true), write to executor disks and are faster, but if an executor is lost the checkpoint goes with it and the job fails, so reserve them for clusters without spot preemption.
Second, hubs: in large-star the minimum of a vertex with ten million neighbours is broadcast to all of them, and the shuffle partition holding that vertex becomes the straggler. Filter hubs semantically first, as in the example, then let the broadcast threshold or AQE's skew-join handling deal with what remains; see Spark data skew for the general techniques. Third, the shuffle partition count: start near two to three times the executor core count and size so each partition holds a few hundred megabytes of edges.
Failure modes
- Job fails asking for a checkpoint directory: set
setCheckpointDirto HDFS or an object store before calling the algorithm, or enable local checkpoints deliberately. - One giant component: a null, empty or default key is joining unrelated vertices. Query the top-degree vertices and add them to the filter.
- Rounds keep climbing on GraphX: the graph has long chains. Switch to
two_phaseorrandomized_contractionrather than raising the iteration cap. - Driver slows down every round: plan growth. Lower the checkpoint interval and confirm checkpoints are actually written.
- Executors lost mid-run with local checkpoints: the checkpoint vanished with the executor. Use reliable storage on preemptible clusters.
- Downstream joins break after a rerun: component IDs changed. Add the overlap-based stable ID mapping.
- Directed results expected: connected components ignores direction. If direction matters, use strongly connected components, which is a different and more expensive algorithm.
Trade-offs
Choose two_phase as the default: deterministic, well-tested, and its minimum-ID labels are at least reproducible when the graph does not change. Choose randomized_contraction when memory is the constraint or benchmarks on your own data favour it, and budget for the stable-ID mapping because its labels are random. Keep GraphX for small-diameter graphs already in RDD pipelines. Consider a hybrid when most components are tiny: run union-find inside each partition with mapPartitions to collapse local structure, then run a distributed algorithm only on the much smaller graph of partition representatives. It cuts rounds and shuffle volume, at the cost of code you must maintain and test yourself. And remember the cheapest optimisation is the hub filter: it improves both run time and the answer.
What to do next
- Measure your graph first: vertex count, edge count, degree distribution and the top 100 vertices by degree.
- Model many-to-many links as a bipartite graph with typed, normalised key vertices.
- Filter null, empty and default keys and set a degree cap justified by your domain.
- Set a reliable checkpoint directory and run
two_phasewith a checkpoint interval of 2. - Plot the component size distribution and investigate any component above an agreed size.
- Log rounds, edges per round and run time, and alert when any of them jumps.
- Add an overlap-based stable ID mapping before any downstream table depends on component IDs.
- Benchmark
randomized_contractionand AQE mode on a copy of production data before switching.