Pregel Tutorial

This tutorial covers GraphFrames' graphframes.lib.Pregel API for developing scalable, iterative graph algorithms using Apache Spark 4.x. We will implement progressively complex algorithms — from simple degree counting to path-tracing algorithms — using the same Stack Exchange knowledge graph from the Motif Finding Tutorial.

Pregel is a vertex-centric programming model for distributed graph processing. It was introduced by Google engineers in 2010 and has become the foundation for graph computation at scale. GraphFrames implements the Pregel model using Apache Spark DataFrames, giving you the full power of Spark's query optimizer and distributed execution engine behind a clean, declarative API.

By the end of this tutorial, you will understand how to think in Pregel — how to decompose graph problems into local vertex computations that converge to global solutions through iterative message passing. This is a fundamentally different way of thinking about graph algorithms, and it unlocks computations that are difficult or impossible to express with traditional graph queries or even GraphFrames' built-in algorithms.

The complete source code for this tutorial is in .

Prerequisites

Before starting this tutorial, ensure you have:

For this tutorial, you'll need GraphFrames version 0.11.0 or later. Check your version:

import graphframes
print(graphframes.__version__)

What is Pregel?

Pregel is a bulk synchronous parallel (BSP) system for large-scale graph processing described in the landmark 2010 paper Pregel: A System for Large-Scale Graph Processing from Malewicz et al. at Google.

The core idea is deceptively simple: think like a vertex. Instead of writing an algorithm that operates on the entire graph at once, you write a function that executes at each vertex independently. The function can read incoming messages, update the vertex's state, and send messages to neighboring vertices. The system handles distribution, synchronization, and fault tolerance.

Pregel BSP Compute Dataflow
CME 323: Distributed Algorithms and Optimization, Reza Zadeh, Databricks and Stanford

Computation proceeds in a series of supersteps. In each superstep:

  1. Compute: Every active vertex executes its function, reading messages from the previous superstep
  2. Communicate: Vertices send messages along their edges to neighbors
  3. Barrier: The system synchronizes — no vertex proceeds to the next superstep until all vertices have finished the current one
Bulk Synchronous Parallel Model
The BSP model: Compute → Communicate → Barrier, repeated until convergence

This barrier synchronization is what makes Pregel algorithms easy to reason about. At any point during execution, you know that all vertices are in the same superstep. There are no race conditions, no stale reads, no distributed coordination headaches. You trade some potential parallelism for massive simplification of the programming model.

Pregel Vertex State Machine
Vertex state machine from the Pregel: A System for Large-Scale Graph Processing: vertices alternate between active and inactive states

Vertices can vote to halt — marking themselves inactive. An inactive vertex is woken up when it receives a new message. When all vertices have voted to halt and there are no messages in transit, the algorithm terminates. This is how Pregel algorithms converge: vertices stop updating when their state stabilizes.

Pregel supersteps from the original paper
Superstep progression from the Pregel: A System for Large-Scale Graph Processing: vertices send messages and receive them in the next superstep

The Power of Local Computation

Pregel's "think like a vertex" paradigm is a profound shift from how most engineers approach graph problems. When you sit down with a graph database and write a Cypher or Gremlin query, you're thinking globally — "find all paths from A to B," "count triangles in the graph," "return the top-10 most central nodes." These are global questions that require the system to traverse large portions of the graph.

In Pregel, you think locally. Your vertex function sees only:

From this limited local view, global solutions emerge through iteration. PageRank doesn't require any vertex to know the entire graph topology — each vertex just needs to know its own PageRank value and out-degree. Through repeated rounds of message passing, the globally correct PageRank values emerge from purely local computations.

This local-to-global property is what makes Pregel algorithms scalable. Each vertex's computation is independent and can run in parallel. The only synchronization point is the barrier between supersteps. This is fundamentally different from graph databases, which often require global locking or distributed transactions for complex queries.

At its heart, Pregel is bulk synchronous parallel processing applied to graphs — or as Dean and Ghemawat (2004) might put it, it is MapReduce adapted for the structure of graph data, where the "map" operation is vertex computation and the "reduce" operation is message aggregation.

Why Pregel?

You might wonder: when do I need Pregel instead of GraphFrames' built-in algorithms like pageRank() or connectedComponents()? Well, in fact both of these are implemented using Pregel :) Pregel is for when the built-in algorithms are not enough. It is a general-purpose framework for writing any iterative graph algorithm.

Pregel excels at problems that require:

Pregel is not the right tool for:

The GraphFrames Pregel API

GraphFrames implements Pregel as a DataFrame-based API with a builder pattern. Here is the general structure:

from graphframes.lib import Pregel

result = graph.pregel \
    .setMaxIter(n) \
    .withVertexColumn("state", initial_expr, update_expr) \
    .sendMsgToDst(message_expr) \
    .aggMsgs(aggregation_expr) \
    .run()

Each part maps directly to the Pregel model:

Method Purpose Pregel Concept
setMaxIter(n) Maximum supersteps Termination condition
withVertexColumn(name, init, update) Define vertex state Vertex function
sendMsgToDst(expr) / sendMsgToSrc(expr) Define messages Edge communication
aggMsgs(expr) Aggregate messages per vertex Message combining
run() Execute the algorithm Superstep loop

The key insight is that everything is expressed as Spark SQL column expressions. You don't write Python loops or callbacks — you describe what to compute, and Spark's optimizer figures out how to execute it efficiently across your cluster.

For full API details, see the graphframes.lib.Pregel Python documentation and the org.graphframes.lib.Pregel Scala documentation.

How GraphFrames Implements Pregel Under the Hood

It helps to understand what happens when you call .run(). The implementation in org.graphframes.lib.Pregel translates the Pregel model into DataFrame operations:

  1. Initialize: The vertices DataFrame is expanded with the initial values of all withVertexColumn definitions, plus an internal active flag column.

  2. Build triplets: For each iteration, GraphFrames constructs triplets — DataFrames where each row contains a source vertex, an edge, and a destination vertex. These triplets are what make Pregel.src(), Pregel.dst(), and Pregel.edge() references work in your message expressions.

  3. Generate messages: The message expressions (sendMsgToDst, sendMsgToSrc) are evaluated over the triplets, producing a DataFrame of (targetvertexid, message) pairs. Null messages are automatically filtered out.

  4. Aggregate: Messages are grouped by target vertex ID and aggregated using the aggMsgs expression.

  5. Update: The aggregated messages are LEFT-OUTER joined with the current vertices. The update expressions in withVertexColumn are evaluated, producing the new vertex state. Vertices that received no messages see null for Pregel.msg().

  6. Checkpoint: Every checkpointInterval iterations (default: 2), the vertex DataFrame is checkpointed to prevent the Spark query plan from growing linearly with the number of iterations. Without this, Spark would build a plan with hundreds of joins for a 50-iteration algorithm, which would crash the driver.

  7. Termination check: If early stopping or vertex voting is enabled, a Spark action is triggered to check the condition.

Automatic optimization: GraphFrames analyzes your message expressions to determine if the destination vertex state is actually needed. If your messages only reference Pregel.src() and Pregel.edge() columns (not Pregel.dst()), the implementation skips the second join entirely — a significant performance optimization for algorithms like PageRank.

Understanding this implementation helps you write better Pregel algorithms. For example:

Data Setup

⚠️ Before continuing: Complete the Data Setup Tutorial first. That tutorial covers downloading the Stack Exchange dataset, converting it to Parquet, and loading nodes_df, edges_df, and g. The examples below assume those objects already exist.

The Stack Exchange graph — specifically the stats.meta.stackexchange.com dataset — contains ~130K nodes (Users, Questions, Answers, Votes, Badges, Tags, PostLinks) and ~97K edges across 8 relationship types. It is small enough to run on a laptop but complex enough to demonstrate real graph algorithms.

You will also need these imports for the Pregel examples:

import pyspark.sql.functions as F
from graphframes import GraphFrame
from graphframes.lib import Pregel


# Must set checkpoints folder to use Pregel API
spark.sparkContext.setCheckpointDir("/tmp/spark-checkpoints")

Example 1: Degree Centrality with AggregateMessages

Before diving into Pregel, let's build intuition with a simpler API: graphframes.GraphFrame.aggregateMessages. AggregateMessages performs a single pass of message sending and aggregation — no iterations, no evolving state. It is the right tool for simple, one-shot graph aggregations.

The most basic graph metric is in-degree: how many edges point to each vertex. In our Stack Exchange graph, a User's in-degree tells us how many badges they've earned, a Question's in-degree tells us how many answers and votes it has received, and so on.

In-degree computation with AggregateMessages
AggregateMessages: each source sends 1 to its destination, destinations sum their messages
from graphframes.lib import AggregateMessages as AM

am_in_degrees = g.aggregateMessages(
    F.count(AM.msg).alias("in_degree"),
    sendToDst=F.lit(1),
)

There is a subtlety here: vertices with no incoming edges (in-degree = 0) won't appear in the results at all, because they received no messages. We fix this with a LEFT JOIN:

complete_in_deg = (
    g.vertices.select("id", "Type")
    .join(am_in_degrees, on="id", how="left")
    .na.fill(0, ["in_degree"])
)

complete_in_deg.groupBy("in_degree").count().orderBy("in_degree").show(10)
+---------+-----+
|in_degree|count|
+---------+-----+
|        0|81735|
|        1|43165|
|        2|  341|
|        3|  218|
|        4|  289|
|        5|  326|
|        6|  371|
|        7|  318|
|        8|  338|
|        9|  304|
+---------+-----+

Most nodes have zero in-degree (they are source-only nodes like Votes that cast votes but don't receive edges). The distribution follows a power law, which is typical for real-world networks. The power law means that a few nodes have very high in-degree (popular questions with hundreds of votes, prolific users with thousands of badges) while the vast majority have low in-degree. Note that a power law distribution with supernodes can kill Pregel performance since each degree is a new message.

A degree-by-degree table only shows the first few rows of a distribution that spans several orders of magnitude — a poor fit for a linear histogram. Instead, we bucket in-degrees into powers of two (0, 1, 2–3, 4–7, 8–15, ...) and draw the bars on a log scale:

import math


def log_hist(df: DataFrame) -> None:
    """Print a text histogram, with bar length on a log scale"""

    # Bucket the in-degrees into powers of two: 0, 1, 2-3, 4-7, 8-15, ...
    histogram = (
        df
        .withColumn(
            "bucket",
            F.when(F.col("in_degree") == 0, F.lit(-1))
            .otherwise(F.floor(F.log2("in_degree")))
        )
        .groupBy("bucket")
        .agg(F.count("*").alias("num_vertices"))
        .orderBy("bucket")
        .collect()
    )

    print(f"{'in_degree':>9}  {'vertices':>9}")
    for row in histogram:
        if row["bucket"] == -1:
            label = "0"
        else:
            low = 2 ** row["bucket"]
            high = 2 ** (row["bucket"] + 1) - 1
            label = str(low) if low == high else f"{low}-{high}"
        bar = "#" * max(1, round(10 * math.log10(max(row["num_vertices"], 1))))
        print(f"{label:>9}  {row['num_vertices']:>9}  {bar}")


log_hist(complete_in_deg)
in_degree   vertices
        0      81735  #################################################
        1      43165  ##############################################
      2-3        559  ###########################
      4-7       1304  ###############################
     8-15       2004  #################################
    16-31        855  #############################
    32-63        112  ####################
   64-127         14  ###########
  128-255          3  #####

The shape of the network is now visible at a glance: the vast majority of vertices receive zero or one edge, and each successive power-of-two bucket thins out until only 3 vertices have an in-degree above 128. Note the pattern in the code: the aggregation happens in Spark, and only the handful of bucket rows are collect()ed to the driver for formatting — a good habit for any summary visualization of large data.

Let's look at who has the highest in-degree:

complete_in_deg.orderBy(F.desc("in_degree")).show(10)

In our Stack Exchange graph, the highest in-degree nodes are typically Questions that have attracted many votes, answers, and tags. This is exactly what we'd expect — Questions are the focal points of a Q&A network. They accumulate relationships from answers, votes, tags, and links.

AggregateMessages is clean and efficient for this kind of single-pass aggregation. But what if we need to compute something that requires multiple iterations — where vertex state evolves based on neighbor states over time? That's where Pregel comes in.

Example 2: Degree Centrality with Pregel

Let's compute the same in-degree metric using Pregel. This is intentionally over-engineered for a single-pass problem, but it gives us a minimal working example to learn the API before we tackle algorithms that need multiple iterations.

In-degree computation with Pregel
Pregel in-degree: initialize to 0, send 1 along each edge, sum at destination
from graphframes.lib import Pregel


pregel_in_degree = g.pregel \
    .setMaxIter(1) \
    .withVertexColumn(
        "in_degree",
        F.lit(0),                              # Initial value: start at 0
        F.coalesce(Pregel.msg(), F.lit(0))     # Update: use message or keep 0
    ) \
    .sendMsgToDst(F.lit(1)) \
    .aggMsgs(F.sum(Pregel.msg())) \
    .run()

log_hist(pregel_in_degree)

Let's break down each part:

setMaxIter(1): Run for exactly one superstep. Degree only needs a single round of message passing.

withVertexColumn("in_degree", F.lit(0), F.coalesce(Pregel.msg(), F.lit(0))): This defines a new vertex column called in_degree. The second argument (F.lit(0)) is the initial value — every vertex starts with degree 0. The third argument is the update expression — after messages are aggregated, set in_degree to the aggregated message value, or 0 if no messages arrived. Pregel.msg() references the aggregated message column.

sendMsgToDst(F.lit(1)): For every edge in the graph, send the integer value 1 from source to destination. This is the message expression — it can reference source vertex columns with Pregel.src("col"), destination vertex columns with Pregel.dst("col"), and edge columns with Pregel.edge("col").

aggMsgs(F.sum(Pregel.msg())): For each vertex, sum all received messages. Pregel.msg() here references the individual (unaggregated) message column. You can use any Spark aggregation function: sum, min, max, collect_list, etc.

run(): Execute the algorithm and return a DataFrame with the original vertex columns plus the new in_degree column.

Notice the key difference from AggregateMessages: Pregel automatically handles zero-degree vertices. Because withVertexColumn initializes every vertex to 0 and the update expression uses F.coalesce(Pregel.msg(), F.lit(0)), vertices that receive no messages keep their initial value of 0. No LEFT JOIN needed.

AggregateMessages vs. Pregel

Feature AggregateMessages Pregel
Iterations Single pass only Multiple with setMaxIter()
Vertex State Manual (pre-create columns) Automatic (withVertexColumn)
Zero-degree nodes Dropped (need LEFT JOIN) Handled via initial values
Syntax Functional, lower-level Builder pattern, declarative
Best For One-off aggregations Iterative algorithms

For simple metrics like degree, either API works. For everything that follows, Pregel is the only option.

Example 3: PageRank with Pregel

PageRank is the algorithm that launched Google. Defined by Larry Page and Sergey Brin in their 1998 paper The PageRank Citation Ranking: Bringing Order to the Web, it computes the "importance" of each node based on the importance of nodes linking to it. The key insight: a node is important if important nodes point to it.

Simplified PageRank Calculation
A simplified PageRank calculation, from the The PageRank Citation Ranking: Bringing Order to the Web

The PageRank formula for a vertex v is:

PR(v) = (1 - d) / N + d × Σ(PR(u) / out_degree(u))

where d is the damping factor (typically 0.85), N is the total number of vertices, and the sum is over all vertices u that have an edge pointing to v.

This is a natural fit for Pregel:

PageRank iterations showing convergence
PageRank values evolving over iterations: each vertex sends PR/out_degree along its edges

First, we need to compute out-degrees (each vertex needs to know how many outgoing edges it has):

out_degrees = g.outDegrees.withColumnRenamed("outDegree", "out_degree")
pr_vertices = (
    nodes_df.join(out_degrees, on="id", how="left")
    .na.fill(1, ["out_degree"])
)
g_pr = GraphFrame(pr_vertices, edges_df)

num_vertices = g_pr.vertices.count()
damping = 0.85
max_iter = 10

Now the Pregel PageRank implementation:

pr_results = g_pr.pregel \
    .setMaxIter(max_iter) \
    .withVertexColumn(
        "pagerank",
        F.lit(1.0 / num_vertices),
        F.coalesce(Pregel.msg(), F.lit(0.0)) * F.lit(damping)
            + F.lit((1.0 - damping) / num_vertices)
    ) \
    .sendMsgToDst(
        Pregel.src("pagerank") / Pregel.src("out_degree")
    ) \
    .aggMsgs(F.sum(Pregel.msg())) \
    .run()

Initialization: Every vertex starts with 1/N — uniform initial importance.

Message: Each vertex sends pagerank / out_degree to each of its destination vertices. This splits the vertex's importance equally among its outgoing edges. Pregel.src("pagerank") references the source vertex's current pagerank column, and Pregel.src("out_degree") references its out-degree.

Aggregation: Sum all incoming PageRank contributions.

Update: Apply the damping formula. The (1 - d) / N term ensures that every vertex receives a small base amount of PageRank, preventing rank sinks.

Let's see the top Questions and Users by PageRank:

pr_results.filter(F.col("Type") == "Question") \
    .select("id", "Title", "pagerank") \
    .orderBy(F.desc("pagerank")) \
    .show(10, truncate=80)

The top-ranked nodes in a knowledge graph reveal its structure. In Stack Exchange, highly-ranked Users are those who receive many votes on their answers, which in turn point to highly-ranked Questions. Tags that are applied to important questions accumulate PageRank transitively. This is the recursive nature of PageRank — importance flows through the graph topology.

Experiment with different max_iter values. You'll find that PageRank converges quickly — most of the "work" happens in the first 5-6 iterations, with diminishing changes afterward. This is because the damping factor causes the influence of distant nodes to decay exponentially with distance.

Comparing with Built-in PageRank

GraphFrames provides a built-in pageRank() method. Let's compare the rankings — the absolute PageRank values may differ due to normalization, but the relative ordering of vertices should be very similar:

from pyspark.sql import Window

builtin_pr = g.pageRank(resetProbability=1 - damping, maxIter=max_iter)

# Compare rankings rather than absolute values
pregel_ranked = (
    pr_results.select("id", F.col("pagerank").alias("pregel_pr"))
    .withColumn("pregel_rank", F.dense_rank().over(
        Window.orderBy(F.desc("pregel_pr"))
    ))
)
builtin_ranked = (
    builtin_pr.vertices.select("id", F.col("pagerank").alias("builtin_pr"))
    .withColumn("builtin_rank", F.dense_rank().over(
        Window.orderBy(F.desc("builtin_pr"))
    ))
)
comparison = pregel_ranked.join(builtin_ranked, on="id")
rank_corr = comparison.stat.corr("pregel_rank", "builtin_rank")
print(f"Rank correlation: {rank_corr:.4f}")

The rank correlation should be very high (close to 1.0), confirming that both implementations agree on which nodes are important even if the absolute PageRank values differ. The built-in PageRank may use different normalization internally, but the ranking — which is what matters in practice — is essentially the same.

The built-in PageRank also supports convergence tolerance via tol parameter — it stops early when the maximum change in any vertex's PageRank drops below the tolerance. Our Pregel implementation uses a fixed iteration count, but you could add tolerance-based convergence using the vertex voting mechanism (see the Advanced Topics section).

Example 4: Connected Components with Pregel

Connected components identifies groups of vertices that are reachable from each other. It is one of the most fundamental graph algorithms — a building block for community detection, data quality checks, and graph decomposition.

The algorithm is elegantly simple in Pregel:

  1. Each vertex starts with its own ID as its component label
  2. Each vertex sends its current label to all neighbors (both directions)
  3. Each vertex adopts the minimum label it receives
  4. Repeat until no labels change

After convergence, all vertices in the same connected component will have the same label — the minimum ID among all vertices in that component.

Connected components label propagation across supersteps
Minimum-label connected components: each vertex starts with its own ID; the component minimum advances one hop per superstep until every vertex in a component shares the same label
cc_vertices = g.vertices.select("id")
cc_graph = GraphFrame(cc_vertices, g.edges.select("src", "dst"))

cc_results = cc_graph.pregel \
    .setMaxIter(20) \
    .setEarlyStopping(True) \
    .withVertexColumn(
        "component",
        F.col("id"),
        F.least(
            F.col("component"),
            F.coalesce(Pregel.msg(), F.col("component"))
        )
    ) \
    .sendMsgToDst(Pregel.src("component")) \
    .sendMsgToSrc(Pregel.dst("component")) \
    .aggMsgs(F.min(Pregel.msg())) \
    .run()

Several important patterns here:

Bidirectional messaging: We use both sendMsgToDst() and sendMsgToSrc(). Directed edges in our graph become undirected for connected components — if A can reach B, then B can reach A. This is a common pattern for algorithms that treat the graph as undirected.

Early stopping: setEarlyStopping(True) tells Pregel to halt when no non-null messages are produced. Once all labels have converged, no vertex sends messages with a smaller label, and the algorithm terminates before reaching maxIter. This is a significant optimization — for many graphs, connected components converges in far fewer than 20 iterations.

Monotonic convergence: The F.least() update ensures labels can only decrease. This guarantees convergence: in the worst case, the minimum ID in each component propagates to every vertex in that component, and no vertex can ever "go backward" to a larger label.

num_components = cc_results.select("component").distinct().count()
print(f"Number of connected components: {num_components:,}")

cc_results.groupBy("component").count() \
    .orderBy(F.desc("count")) \
    .show(10)

In the Stack Exchange graph, you'll typically find one giant component (most nodes are connected through the question-answer-user chain) plus many small components (isolated votes, orphaned badges, etc.).

How Fast Does It Converge?

The convergence speed of connected components depends on the diameter of the graph — the longest shortest path between any two connected vertices. In the worst case (a linear chain of N vertices), it takes N-1 supersteps. In practice, real-world graphs have small diameters due to the small-world property, so convergence is fast.

The early stopping optimization is crucial here. Without it, Pregel would run all 20 iterations even after convergence. With it, the algorithm halts as soon as the minimum labels stop propagating — which often happens in 5-10 iterations for real-world social graphs.

You can verify convergence speed by adding some logging. Run with a few different maxIter values and compare the results — you'll find the same component labels regardless of whether you set maxIter to 10, 20, or 100, as long as it's high enough.

Choosing a Termination Condition

Connected components converges naturally, but "converged" and "stopped" are two different things — Pregel keeps running until you tell it to stop. The Pregel API Reference documents three termination conditions, and picking the wrong one either wastes supersteps or costs more than it saves:

Condition How to set it Stops when
Iteration count setMaxIter(n) n supersteps have run
No new messages setEarlyStopping(True) Every message produced in a superstep is null
Vertex voting setStopIfAllNonActiveVertices(True) plus setInitialActiveVertexExpression(expr) and setUpdateActiveVertexExpression(expr) Every vertex has marked itself inactive

The two dynamic conditions are not free. Both setEarlyStopping and setStopIfAllNonActiveVertices have to look at the data to decide whether to continue, and that means an extra Spark action on every superstep — a full pass over the messages or the vertices, on top of the work the superstep already did. setMaxIter costs nothing because it needs no information about the data.

That trade-off is what makes the choice algorithm-specific:

A practical default: start with setMaxIter(n) alone and a generous n, confirm your algorithm produces stable results, and only then add a dynamic condition — and only if you measured that it actually saves supersteps.

A Note on the Built-in Algorithm

GraphFrames provides g.connectedComponents() which uses a more sophisticated implementation with checkpointing and GraphX integration. For production workloads, the built-in version may be faster due to its optimizations. But the Pregel version is instructive and can be customized — for example, you could add type-aware component detection where edges of certain types are ignored.

Example 5: Shortest Paths with Pregel

Single-source shortest paths computes the minimum number of hops from a source vertex to every other vertex in the graph. In our Stack Exchange graph, this tells us the network distance between entities: how many relationship hops separate a given question from any user, answer, tag, or vote in the network.

The algorithm is a Pregel adaptation of breadth-first search:

  1. The source vertex starts with distance 0; all others start at infinity
  2. Each vertex with a finite distance sends distance + 1 to its neighbors
  3. Each vertex takes the minimum of its current distance and incoming messages
  4. Repeat until no distances improve
Shortest paths propagation
Shortest paths from source A: distances propagate outward, one hop per superstep
# Pick the most-viewed question as our source
popular_question = (
    nodes_df.filter(F.col("Type") == "Question")
    .orderBy(F.desc("ViewCount"))
    .select("id", "Title")
    .first()
)
source_id = popular_question["id"]

sp_vertices = g.vertices.select("id")
sp_graph = GraphFrame(sp_vertices, g.edges.select("src", "dst"))

INF = 999999  # Represents unreachable

sp_results = sp_graph.pregel \
    .setMaxIter(10) \
    .setEarlyStopping(True) \
    .withVertexColumn(
        "distance",
        F.when(F.col("id") == source_id, F.lit(0)).otherwise(F.lit(INF)),
        F.least(
            F.col("distance"),
            F.coalesce(Pregel.msg(), F.lit(INF))
        )
    ) \
    .sendMsgToDst(
        F.when(
            Pregel.src("distance") < F.lit(INF),
            Pregel.src("distance") + F.lit(1)
        )
    ) \
    .sendMsgToSrc(
        F.when(
            Pregel.dst("distance") < F.lit(INF),
            Pregel.dst("distance") + F.lit(1)
        )
    ) \
    .aggMsgs(F.min(Pregel.msg())) \
    .run()

Key patterns:

Conditional messaging: The F.when(Pregel.src("distance") < F.lit(INF), ...) guard ensures that only vertices with finite distances send messages. Vertices that haven't been reached yet don't waste computation sending infinity + 1. null messages are automatically filtered by Pregel.

Bidirectional again: We send messages in both directions to treat the graph as undirected for reachability. If you want directed shortest paths, use only sendMsgToDst.

BFS-like expansion: Each superstep extends the frontier by one hop. The distance distribution tells us the structure of the graph: how quickly information can flow through the network.

sp_results.filter(F.col("distance") < INF) \
    .groupBy("distance").count() \
    .orderBy("distance").show(20)

You can see the messages spread in the distance distribution:

+--------+-----+
|distance|count|
+--------+-----+
|       0|    1|
|       1|   69|
|       2|  468|
|       3| 4375|
|       4|24217|
|       5|25126|
|       6| 2160|
|       7|   26|
+--------+-----+

The distance distribution reveals the small-world property of the Stack Exchange graph. Most vertices are reachable within a handful of hops.

Understanding the Wavefront

Shortest paths in Pregel is essentially BFS implemented as message passing. The algorithm creates an expanding wavefront: at superstep 1, all direct neighbors of the source get distance 1. At superstep 2, their neighbors get distance 2. And so on, one hop per superstep.

BFS wavefront expanding one hop per superstep
The shortest-path wavefront: only the source starts at distance 0; each superstep reaches the next ring of neighbors while unreached vertices remain INF

This wavefront property is why the distance distribution is so informative. Each row in the distribution corresponds to one superstep's work. If distance 3 has the most vertices, it means the third ring of neighbors from the source is the densest — which tells you something about the graph's local structure around your source vertex.

In our Stack Exchange graph, the source question connects directly (distance 1) to its answers, the users who asked it, tags applied to it, and votes cast on it. At distance 2, we reach the badges those users have, other questions they've asked, other answers to the same question, and so on. By distance 3-4, we've typically reached most of the connected graph.

The INF value we use (999999) is a practical choice. A more mathematically pure implementation would use None/null, but nulls in Spark column expressions require more careful handling with coalesce. The large integer approach is simpler and works correctly as long as your graph has fewer than 999,999 vertices in any path — which is true for virtually all real-world graphs.

Advanced: Memory Optimization

For large graphs with many vertex columns, Pregel constructs triplets (source vertex + edge + destination vertex) that can be memory-intensive. GraphFrames provides required_src_columns() and required_dst_columns() to specify which columns are actually needed:

sp_graph.pregel \
    .setMaxIter(10) \
    .withVertexColumn("distance", ..., ...) \
    .sendMsgToDst(Pregel.src("distance") + F.lit(1)) \
    .required_src_columns("distance") \
    .required_dst_columns("distance") \
    .aggMsgs(F.min(Pregel.msg())) \
    .run()

The id column and active flag are always included automatically. This optimization can dramatically reduce memory usage for graphs with wide vertex schemas (like our Stack Exchange graph with 50+ columns).

Example 6: Reputation Propagation

Now we get to the real power of Pregel: algorithms that don't exist as built-in functions. This is where vertex-centric thinking pays off.

In Stack Exchange, users have reputation scores that reflect their trustworthiness and expertise. But reputation is an attribute of users, not of the questions they answer. What if we want to know which questions have been answered by the most reputable users? Which questions have attracted the most expert attention?

This is a reputation propagation algorithm — a form of trust propagation through a network. The concept has roots in the EigenTrust algorithm (Kamvar et al., 2003) for peer-to-peer networks and in trust propagation research for social networks (Chakraborty and Karform, 2012). The core idea is that trust or authority flows through edges, accumulating at destination nodes.

Reputation propagation from users through answers to questions
Reputation flows from Users through Answers to Questions over 2 Pregel supersteps

Here's why this requires Pregel: the reputation must flow two hops — from User to Answer to Question. A single-pass AggregateMessages can only propagate information one hop. You could chain two AggregateMessages calls manually, but Pregel handles multi-hop propagation natively with setMaxIter(2).

We build a subgraph focused on the reputation flow path:

# User nodes with their reputation scores
user_nodes = nodes_df.filter(F.col("Type") == "User").select(
    "id",
    F.col("Reputation").cast("double").alias("reputation"),
    F.lit("User").alias("Type")
)

# Answer nodes with their scores
answer_nodes = nodes_df.filter(F.col("Type") == "Answer").select(
    "id",
    F.col("Score").cast("double").alias("score"),
    F.lit("Answer").alias("Type")
)

# Question nodes
question_nodes = nodes_df.filter(F.col("Type") == "Question").select(
    "id",
    F.col("ViewCount").cast("double").alias("views"),
    F.lit("Question").alias("Type")
)

# Build unified vertex DataFrame
rep_vertices = (
    user_nodes.withColumn("score", F.lit(0.0)).withColumn("views", F.lit(0.0))
    .unionByName(
        answer_nodes.withColumn("reputation", F.lit(0.0))
        .withColumn("views", F.lit(0.0))
    )
    .unionByName(
        question_nodes.withColumn("reputation", F.lit(0.0))
        .withColumn("score", F.lit(0.0))
    )
    .na.fill(0.0)
)

# Edges: User->Answer (Posts) and Answer->Question (Answers)
posts_edges = edges_df.filter(F.col("relationship") == "Posts").select("src", "dst")
answers_edges = edges_df.filter(F.col("relationship") == "Answers").select("src", "dst")
rep_edges = posts_edges.unionByName(answers_edges)

rep_graph = GraphFrame(rep_vertices, rep_edges)

Now the Pregel propagation:

rep_results = rep_graph.pregel \
    .setMaxIter(2) \
    .withVertexColumn(
        "authority",
        F.col("reputation"),  # Users start with their rep; others start at 0
        F.coalesce(Pregel.msg(), F.lit(0.0)) + F.col("authority")
    ) \
    .sendMsgToDst(
        F.when(
            Pregel.src("authority") > F.lit(0),
            Pregel.src("authority")
        )
    ) \
    .aggMsgs(F.sum(Pregel.msg())) \
    .run()

Iteration 1: Users send their reputation scores to the Answers they posted. Each Answer accumulates the reputation of its author.

Iteration 2: Answers (now carrying their author's reputation) send their authority to the Questions they answer. Each Question accumulates the total reputation of all its answerers.

The result is an "authority" score for each Question that reflects the collective expertise of its answer pool. This is a metric you cannot compute with a simple SQL join — it requires understanding the graph topology and propagating values through it.

To understand the subtlety here, consider what happens if you try to compute this without Pregel. You could join Users with Answers (via the "Posts" relationship) and then join Answers with Questions (via the "Answers" relationship), summing reputation along the way. That's a two-hop join, which is possible in SQL. But what if you wanted reputation to flow three hops — through answers to linked questions? Or four hops? Each additional hop requires another join, and the query plan grows quadratically in complexity. Pregel handles arbitrary hop counts by simply increasing setMaxIter.

More importantly, Pregel allows the propagation to be stateful. In our example, each Answer accumulates its author's reputation before forwarding it. With joins, you'd need to pre-compute intermediate results and manage them manually. Pregel's vertex state makes this natural.

This is the key pattern that makes Pregel invaluable for custom graph algorithms: multi-hop, stateful information propagation. Graph databases can do BFS. SQL can do joins. But neither handles iterative, stateful propagation as cleanly as Pregel.

top_questions = (
    rep_results.filter(F.col("Type") == "Question")
    .select("id", "authority")
    .join(
        nodes_df.filter(F.col("Type") == "Question")
        .select("id", "Title", "ViewCount"),
        on="id"
    )
    .select("authority", "Title", "ViewCount")
    .orderBy(F.desc("authority"))
)
top_questions.show(10, truncate=80)
+---------+----------------------------------------------------------------+---------+
|authority|                                                           Title|ViewCount|
+---------+----------------------------------------------------------------+---------+
|3024522.0|         Top $k$ List of Reasons to Close a Question Immediately|     3306|
|1818772.0|                                          Help spread the wealth|     1725|
|1728975.0|                        Internet Support for Statistics Software|    14545|
|1129557.0|                                    Library of helpful responses|      195|
| 992347.0|        Project Reduplication of Deduplication - Cross Validated|      378|
| 960856.0|                                  Community Promotion Ads - 2015|      611|
| 874877.0|Why is a comment that maligns my field not inappropriate for CV?|      719|
| 851696.0|               2017 moderator election Q&A - question collection|      345|
| 777214.0|                                Rich get richer phenomenon on CV|     1145|
| 757480.0|                         Putting Cross Validated stats on resume|     1008|
+---------+----------------------------------------------------------------+---------+
top_questions.stat.corr("authority", "ViewCount")
0.070385635536485

The highest-authority questions are those answered by the most reputable community members — which often (but not always - only a 0.07 Pearson's correlation) corresponds with a significant view count. Divergences between view count and authority reveal questions that are popular but underserved by experts, or niche questions that attracted top-tier answers.

Designing Your Own Pregel Algorithm

Now that we have five algorithms under our belt, let's step back and develop a systematic approach to designing new Pregel algorithms. The goal is to give you a mental framework you can apply to any graph problem you encounter.

The Four-Question Framework

Every Pregel algorithm answers four questions. They are not arbitrary: each one corresponds to a component of the model in the original Pregel paper — the vertex value, the message function, the combiner, and compute() — which is why the GraphFrames API has the shape it does.

Four design questions mapped to Pregel paper concepts and GraphFrames API
Each design question maps to a Pregel paper concept and a GraphFrames API call

Let's walk through how to answer them for a hypothetical new problem: computing the average answer score for each tag in the Stack Exchange graph. Tags connect to Questions, Questions connect to Answers, and Answers have scores. We want each Tag to know the average score of all Answers to its tagged Questions.

Answer scores propagating through Questions to Tags
Worked example: answer scores and counts propagate Answer → Question → Tag; the tag average is total / count

Question 1: What vertex state do I need?

Each Tag needs to accumulate (a) the total score of answers to its questions and (b) the count of such answers. Each vertex gets two columns: total_score and answer_count. Tags initialize both to 0. Answers initialize total_score to their own score and answer_count to 1. Questions and other node types initialize to 0.

Question 2: What messages should neighbors exchange?

The information needs to flow: Answer → Question → Tag. Answers send their (score, count) to the Questions they answer. Questions forward the accumulated (score, count) to their Tags. This is a two-hop propagation, like reputation propagation.

Question 3: How do messages combine?

When a Question receives scores from multiple Answers, it sums both the scores and the counts. When a Tag receives from multiple Questions, it sums again. The aggregation is sum for both fields.

Question 4: How does state update?

After receiving messages, each vertex adds the incoming total to its own total and the incoming count to its own count. The final average is total_score / answer_count, computed after Pregel completes.

Here is the full implementation. We reverse the Tags edges (Tag → Question becomes Question → Tag) so both hops travel with sendMsgToDst, then send a struct message carrying score and count together. Two withVertexColumn definitions unpack the aggregated struct into vertex state — note that withVertexColumn init/update use plain F.col(...) (vertex context), while sendMsgToDst uses Pregel.src(...) (triplet context).

# Subgraph: Tags, Questions, and Answers (Answers carry Score)
tag_nodes = nodes_df.filter(F.col("Type") == "Tag").select(
    "id",
    F.col("TagName").alias("name"),
    F.lit("Tag").alias("Type"),
)
question_nodes = nodes_df.filter(F.col("Type") == "Question").select(
    "id",
    F.lit(None).cast("string").alias("name"),
    F.lit("Question").alias("Type"),
)
answer_nodes = nodes_df.filter(F.col("Type") == "Answer").select(
    "id",
    F.lit(None).cast("string").alias("name"),
    F.col("Score").cast("double").alias("score"),
    F.lit("Answer").alias("Type"),
)

avg_vertices = (
    tag_nodes.withColumn("score", F.lit(0.0))
    .unionByName(question_nodes.withColumn("score", F.lit(0.0)))
    .unionByName(answer_nodes)
    .na.fill({"score": 0.0})
)

avg_vertices.sample(0.002).show()

Creating a single node type is helpful for implementing Pregel algorithms, but you can also use conditionals to implement vertex state updates.

+--------------------+--------+--------+-----+
|                  id|    name|    Type|score|
+--------------------+--------+--------+-----+
|5803e875-3f19-496...|homework|     Tag|  0.0|
|6bfb6d29-5bcf-491...|    NULL|Question|  0.0|
|97914dc4-00c5-4d5...|    NULL|Question|  0.0|
|b6f34e27-dbbf-48a...|    NULL|Question|  0.0|
|3004092b-db46-414...|    NULL|  Answer|  8.0|
|ecea740b-8341-466...|    NULL|  Answer|  7.0|
|81cbf28a-014a-461...|    NULL|  Answer| 19.0|
|c58a303e-8edb-423...|    NULL|  Answer|  7.0|
+--------------------+--------+--------+-----+

Note how often we have created new GraphFrames of a single type as a method of problem solving in this tutorial.

# Answers → Questions (as-is); reverse Tags so Questions → Tags
answers_edges = edges_df.filter(F.col("relationship") == "Answers").select("src", "dst")
questions_to_tags = (
    edges_df.filter(F.col("relationship") == "Tags")
    .select(F.col("dst").alias("src"), F.col("src").alias("dst"))
)
avg_edges = answers_edges.unionByName(questions_to_tags)
avg_graph = GraphFrame(avg_vertices, avg_edges)

# Two-hop propagation: Answer → Question → Tag
avg_results = (
    avg_graph.pregel
    .setMaxIter(2)
    # Q1 + Q4: vertex state columns (init / update use F.col, not Pregel.src)
    .withVertexColumn(
        "total_score",
        F.when(F.col("Type") == "Answer", F.col("score")).otherwise(F.lit(0.0)),
        F.col("total_score")
        + F.coalesce(Pregel.msg().getField("score"), F.lit(0.0)),
    )
    .withVertexColumn(
        "answer_count",
        F.when(F.col("Type") == "Answer", F.lit(1.0)).otherwise(F.lit(0.0)),
        F.col("answer_count")
        + F.coalesce(Pregel.msg().getField("count"), F.lit(0.0)),
    )
    # Q2: messages — triplet context, so sender state is Pregel.src(...)
    .sendMsgToDst(
        F.when(
            Pregel.src("answer_count") > F.lit(0),
            F.struct(
                Pregel.src("total_score").alias("score"),
                Pregel.src("answer_count").alias("count"),
            ),
        )
    )
    # Q3: one aggMsgs expression — struct of sums, not nested aggregates
    .aggMsgs(
        F.struct(
            F.sum(Pregel.msg().getField("score")).alias("score"),
            F.sum(Pregel.msg().getField("count")).alias("count"),
        )
    )
    .run()
)

tag_avg = (
    avg_results.filter(F.col("Type") == "Tag")
    .filter(F.col("answer_count") > 0)
    .withColumn("avg_answer_score", F.col("total_score") / F.col("answer_count"))
    .select("name", "total_score", "answer_count", "avg_answer_score")
    .orderBy(F.desc("avg_answer_score"))
)
tag_avg.show(10, truncate=40)
+-------------------------+-----------+------------+------------------+
|                     name|total_score|answer_count|  avg_answer_score|
+-------------------------+-----------+------------+------------------+
|            stackexchange|       35.0|         1.0|              35.0|
|       math-stackexchange|       28.0|         1.0|              28.0|
|               winterbash|      162.0|         8.0|             20.25|
|                      faq|     1105.0|        66.0|16.742424242424242|
|            participation|     1090.0|        85.0|12.823529411764707|
|datascience-stackexchange|      375.0|        30.0|              12.5|
|            display-names|       25.0|         2.0|              12.5|
|          site-statistics|      576.0|        47.0| 12.25531914893617|
|           community-user|       12.0|         1.0|              12.0|
|                  careers|      192.0|        16.0|              12.0|
+-------------------------+-----------+------------+------------------+

Iteration 1: Answers send (score, 1) to their Questions. Each Question accumulates the total score and answer count for its answers.

Iteration 2: Questions forward those accumulated (total_score, answer_count) structs to their Tags. Each Tag ends up with the aggregate score and count of all answers to questions carrying that tag.

The post-Pregel division total_score / answer_count is the average answer score per tag — a two-hop, stateful aggregation that maps directly from the four questions above. The same decomposition works for any graph problem that involves information propagation.

Common Algorithm Templates

Most Pregel algorithms fall into a few templates:

Propagation (PageRank, reputation): Vertex values flow outward along edges and accumulate. Each vertex updates its state based on the sum (or weighted sum) of incoming values. Converges through damping or fixed iterations.

Label spreading (connected components, community detection): Vertices adopt labels from their neighbors using a consensus rule (minimum, majority, random). Converges when labels stabilize.

Frontier expansion (shortest paths, BFS): A wavefront of "active" vertices expands outward from seed vertices. Each superstep extends the frontier by one hop. Converges when the frontier can't expand further.

Belief propagation (probabilistic inference): Vertices exchange probability distributions and update their beliefs based on Bayesian inference rules. Used in graphical model inference and recommendation systems.

Recognizing which template fits your problem is half the battle. The other half is handling the edge cases: null messages, zero-degree vertices, convergence criteria, and the inevitable off-by-one error in your message expression.

Example 7: Debug Trace — Understanding Pregel's Mechanics

When developing a new Pregel algorithm, it is invaluable to see exactly what messages flow where and how vertex state evolves. This example is purely educational: we track the path of messages through a small test graph to make the message-passing mechanics visible.

Debug trace path accumulation across Pregel supersteps
Message path tracing on A→B, A→C, B→C, C→D: each vertex prepends incoming traces and appends its own id; | joins multiple incoming paths
# Small test graph for clarity
test_vertices = spark.createDataFrame(
    [("A", "Alice"), ("B", "Bob"), ("C", "Charlie"), ("D", "David")],
    ["id", "name"],
)
test_edges = spark.createDataFrame(
    [("A", "B"), ("A", "C"), ("B", "C"), ("C", "D")],
    ["src", "dst"],
)
test_graph = GraphFrame(test_vertices, test_edges)

trace_results = test_graph.pregel \
    .setMaxIter(3) \
    .withVertexColumn(
        "trace",
        F.col("id"),
        F.concat_ws(
            " <- ",
            F.coalesce(Pregel.msg(), F.lit("")),
            F.col("id")
        )
    ) \
    .sendMsgToDst(Pregel.src("trace")) \
    .aggMsgs(
        F.concat_ws(" | ", F.collect_list(Pregel.msg()))
    ) \
    .run()

trace_results.select("id", "name", "trace").orderBy("id").show(truncate=False)

The output shows the accumulated message paths at each vertex:

+---+-------+------------------------+
|id |name   |trace                   |
+---+-------+------------------------+
|A  |Alice  | <- A                   |
|B  |Bob    | <- A <- B              |
|C  |Charlie| <- A |  <- A <- B <- C |
|D  |David  |A <- B |  <- A <- C <- D|
+---+-------+------------------------+

Each vertex's trace column shows who influenced it. Reading right-to-left: X <- Y <- Z means Z's state flowed through Y to reach X. The | separator shows independent paths arriving at the same vertex. Notice how vertex D's trace shows the full chain: A's state propagated through B to C and then to D over three supersteps.

This is a debugging technique. When you're developing a new algorithm and the results don't look right, add a trace column to see exactly what messages each vertex receives and how they aggregate. Once the algorithm works, remove the trace.

Debugging Tips for Custom Algorithms

When your Pregel algorithm produces unexpected results, here is a systematic debugging approach:

  1. Start with a small graph (4-6 vertices) where you can compute the expected results by hand.

  2. Add a trace column like Example 7 to see exactly what messages flow where. String concatenation with concat_ws and collect_list makes message flows visible.

  3. Run one iteration at a time by setting setMaxIter(1), then setMaxIter(2), etc. Compare the intermediate results with your manual calculations.

  4. Check for null handling. The most common bug in Pregel algorithms is incorrect null handling. Remember: Pregel.msg() is null when a vertex receives no messages. Your update expression must handle this case explicitly, usually with F.coalesce(Pregel.msg(), default_value).

  5. Verify message direction. sendMsgToDst sends from source to destination (following the edge direction). sendMsgToSrc sends from destination to source (against the edge direction). For undirected algorithms, you need both.

  6. Check aggregation semantics. F.sum() treats nulls as zero. F.min() ignores nulls. F.collect_list() excludes nulls. Make sure your aggregation function does what you expect.

  7. Watch for zero-degree vertices. Vertices with no incoming messages won't trigger the update expression (they'll see null for Pregel.msg()). Make sure your initial values and update expressions handle this correctly.

These debugging patterns apply to every Pregel algorithm. Internalizing them will save you hours of confusion when developing custom algorithms on real data.

Advanced Topics

Advanced topics include convergence strategies, performance best practicaes and when not to use Pregel.

Convergence Strategies

GraphFrames Pregel provides three ways to control when the algorithm stops:

setMaxIter(n): The simplest — run for exactly n supersteps. Good for algorithms with known iteration counts (PageRank typically converges in 10-20 iterations).

setEarlyStopping(True): Stop when no non-null messages are produced. Good for BFS-like algorithms (shortest paths, connected components) where the frontier eventually stops expanding. Warning: checking for empty messages is a Spark action, so it adds overhead per iteration. Don't use for algorithms like PageRank where messages are never null.

Vertex voting via setStopIfAllNonActiveVertices(True) with setInitialActiveVertexExpression(expr) and setUpdateActiveVertexExpression(expr): The most general approach. Each vertex has an active flag. You define expressions for when vertices should be active. When all vertices become inactive, the algorithm terminates. This is closest to the original Pregel paper's "vote to halt" mechanism.

For PageRank with tolerance-based convergence, vertex voting looks like:

graph.pregel \
    .setMaxIter(100) \
    .setStopIfAllNonActiveVertices(True) \
    .setUpdateActiveVertexExpression(
        F.abs(
            F.col("pagerank")
            - (F.coalesce(Pregel.msg(), F.lit(0.0)) * F.lit(damping)
               + F.lit((1.0 - damping) / num_vertices))
        ) > F.lit(0.001)
    ) \
    .withVertexColumn("pagerank", ..., ...) \
    ...

Note that the active vertex expression must recompute the new PageRank value to compare with the old one — you cannot simply compare F.col("pagerank") with Pregel.msg(), since the message is the raw sum of incoming contributions, not the final damped PageRank.

Performance Best Practices

Start with a small graph. Test your algorithm on a small test graph (like Example 7) before running on the full dataset. This catches logical errors quickly.

Use setEarlyStopping(True) for frontier-based algorithms. Connected components and shortest paths often converge much faster than maxIter. Early stopping avoids wasting computation on empty supersteps.

Use setSkipMessagesFromNonActiveVertices(True) for algorithms with active vertex tracking. This avoids generating messages from vertices whose state hasn't changed, dramatically reducing message volume for sparse-update algorithms like shortest paths.

Checkpoint regularly. Pregel checkpoints every 2 iterations by default (configurable via setCheckpointInterval). For long-running algorithms, this prevents Spark's query plan from growing unboundedly. Set a checkpoint directory via sc.setCheckpointDir("/path").

Consider storage levels. For very large graphs, use setIntermediateStorageLevel(StorageLevel.DISK_ONLY) to avoid OOM errors at the cost of slower iteration.

Use required_src_columns / required_dst_columns when your vertex schema is wide but your message expressions only reference a few columns. This can reduce shuffle data by orders of magnitude.

Unpersist results. Pregel returns a persisted DataFrame. Call .unpersist() when you're done to free memory:

results = graph.pregel.setMaxIter(10)...run()
# ... use results ...
results.unpersist()

When NOT to Use Pregel

Pregel is powerful, but it's not always the right tool:

Conclusion

We covered a lot of ground in this tutorial. Starting from the simplest possible graph metric (in-degree), we progressed through increasingly sophisticated algorithms to custom reputation propagation — an algorithm that doesn't exist in any graph library but was straightforward to implement with Pregel.

Here is the core pattern that every Pregel algorithm follows:

result = graph.pregel \
    .setMaxIter(n) \
    .withVertexColumn("state", initial_value, update_function) \
    .sendMsgToDst(message_expression) \
    .aggMsgs(aggregation_function) \
    .run()

The art of Pregel programming is in choosing the right:

If you can answer these four questions for your problem, you can implement it in Pregel.

Pattern Library

Here is a summary of the algorithm patterns we implemented, along with their key characteristics. Use this as a reference when designing your own Pregel algorithms:

Algorithm Iterations Message Direction Aggregation State Type Convergence
In-Degree 1 dst only sum integer Single pass
PageRank 10-20 dst only sum float Damped iteration
Connected Components variable both (undirected) min string/ID Early stopping
Shortest Paths variable both (undirected) min integer Early stopping
Reputation Propagation 2 dst only sum float Fixed iterations
Debug Trace variable dst only collect_list + concat string Fixed iterations

Common patterns:

Key Takeaways

  1. Think like a vertex. Your algorithm is a function that runs at each vertex, seeing only local information and messages from neighbors. Global behavior emerges from local interactions.

  2. Pregel handles the hard parts. Distribution, synchronization, fault tolerance, and optimization are handled by the framework. You focus on the algorithm logic.

  3. Start with AggregateMessages for single-pass problems. Move to Pregel when you need multiple iterations or evolving vertex state.

  4. Use early stopping for algorithms that converge before maxIter. It's free performance for BFS-like algorithms.

  5. Real-world data reveals real-world structure. Running these algorithms on Stack Exchange data (or your own data) reveals patterns that synthetic examples can't show.

  6. Pregel enables algorithms that don't exist elsewhere. The reputation propagation example demonstrates a custom, domain-specific graph algorithm that would be complex to implement with joins alone. When you encounter a graph problem that no library solves, Pregel gives you the building blocks to solve it yourself.

  7. The four questions of Pregel design. For any graph problem, ask: What vertex state do I need? What messages should neighbors exchange? How do messages combine? How does state update? If you can answer these, you can implement it in Pregel.

Further Reading