Large-scale graph processing
Running algorithms over graphs with billions of edges: the Pregel vertex-centric model, PageRank, connected components, friend-of-friend recommendations, partitioning graphs, offline batch versus online queries, incremental updates, and where graph processing appears in system design.
Reading is half of it. See this used in a real interview: walk through Design a Web Crawler →
Social networks, the web, payment networks and road maps are graphs with billions of nodes and edges. Some questions about them are simple lookups ("who does Alice follow?"), served online from adjacency lists. Others need the whole graph: ranking web pages, finding fraud rings, suggesting people you may know. Those run as large offline graph computations. Knowing the models and the trade-offs helps in crawler, social and fraud designs.
Online queries versus offline processing
- Online: one or two hops from a node, answered in milliseconds from adjacency lists in a key-value store, wide-column store or graph database. See graph data.
- Offline: algorithms touching every node, many times, run in batch over snapshots (hours), with results written to stores the online system reads.
A common pattern: compute expensive graph features offline (PageRank, communities, candidate recommendations) and serve them online by key.
The Pregel model
Google's Pregel introduced vertex-centric computation, used by Apache Giraph, Spark GraphX and others:
- Computation proceeds in supersteps.
- In each superstep, every active vertex reads messages sent to it in the previous step, updates its value, and sends messages along its edges.
- Vertices vote to halt when they have nothing to do; the job ends when all have halted and no messages are in flight.
"Think like a vertex" makes many algorithms simple to express and parallelise across machines, with checkpointing for fault tolerance.
Classic algorithms
PageRank: a page is important if important pages link to it. Each superstep, every page sends its rank divided by its out-degree to the pages it links to; each page's new rank is a damping factor times the sum received, plus a small constant. A few dozen iterations converge. Used for web search ranking and influence scores. See Design a Web Crawler.
Connected components: each vertex starts with its own id as label and repeatedly adopts the smallest label among its neighbours; when nothing changes, vertices sharing a label are connected. Used to find fraud rings linked by shared devices or cards. See fraud detection.
Friends of friends / people you may know: for each user, count how many of their friends know each non-friend; rank by mutual friends plus other signals. Doing this for celebrities with millions of connections explodes, so cap degrees, sample, or skip very high-degree nodes. See Design Instagram.
Shortest paths and centrality: routing (with specialised techniques), degrees of separation, finding key nodes. See routing and shortest paths.
Community detection and label propagation: clusters of densely connected users, useful for recommendations and for spotting coordinated behaviour. See trust and safety.
Partitioning graphs
Graph computations are communication-heavy: messages cross machines whenever edges do.
- Hash partitioning of vertices is simple and balanced but cuts most edges.
- Edge partitioning (vertex-cut) splits high-degree vertices across machines, which helps with power-law graphs where a few nodes have millions of edges.
- Locality-aware partitioning (by community or geography) reduces cross-machine edges, at the cost of balance and preprocessing.
High-degree nodes (celebrities, hub pages) are the main skew problem; handle them specially. See hot keys and skew.
Alternatives to Pregel
- MapReduce or Spark joins: express each iteration as joins of edge and vertex tables; works with existing data infrastructure but writes intermediate results each iteration. See MapReduce and Spark.
- Single-machine graph engines with large memory and compressed graphs can process surprisingly large graphs (billions of edges) faster than clusters, by avoiding network costs.
- Graph databases for interactive multi-hop queries on moderate data sizes, not whole-graph analytics.
Incremental updates
Graphs change constantly. Recomputing everything daily is common and simple. For fresher results, maintain features incrementally (update mutual-friend counts when an edge is added) for the most valuable signals, and recompute the rest in batch. Stream processors can maintain simple graph aggregates in real time. See batch and stream processing.
Serving results
Write outputs (scores, candidate lists, component ids) to a key-value store or feature store keyed by vertex id, versioned per run so you can roll back a bad computation. See feature stores and ML serving.
In the interview
"Follow edges live in a wide-column store for online one-hop queries. Nightly, a Spark GraphX job over the edge snapshot computes people-you-may-know candidates by mutual-friend counts (capping high-degree nodes) and connected components over shared-device edges for fraud; results are written versioned to a key-value store and read by the recommendation and risk services." For a crawler: "PageRank runs as an iterative batch job over the link graph and feeds search ranking."
Checklist
- Online one-hop queries separated from offline whole-graph computation.
- Vertex-centric (Pregel) or join-based iterations, with checkpoints.
- PageRank, components, mutual friends, communities as standard tools.
- Partitioning aware of high-degree nodes and cross-machine edges.
- Daily recomputation plus incremental updates for key features.
- Versioned outputs served by key.