The cut turns hops into RPCs — graph sharding, locality and network round trips
You cannot shard a graph like a key-value store because a graph's value IS the edges between keys, so any partition cuts edges and a traversal crossing the cut replaces each pointer dereference with a network RPC — making a 3-hop walk that zig-zags across the boundary three network round trips instead of three fast follows.
On one primary every pointer was local. Split the graph and the cut severs edges — a hop across it becomes a network RPC, so O(1) becomes O(network). That's unavoidable for a connected graph. But HOW you cut decides how MUCH crosses — and on real, power-law graphs, the obvious cut (assign each vertex to a machine) lands every one of a celebrity's edges on a single node.
Scene 11
The cut turns hops into RPCs
- Watch
- Try it
- Predict
- Capture
On one machine, a friends-of-friends walk from Alice was three pointer dereferences — microseconds. Now watch what splitting the graph across two machines does. First the slice drops and the edges it severs flash. That severed slice is the partition cut. Then the same 3-hop walk runs — but now it zig-zags across the cut, and each crossing is no longer a pointer follow. It's a network round trip: a cross-partition hop, an RPC. The hop count is identical; the cost is a thousand times higher.
Highlighted lines are the ones running in the diagram right now.
def follow_edge(cur, edge):if machine_of(edge.other_end) == this_machine:# local hop: index-free adjacency, ~1 µsreturn deref_pointer(edge) # pointer followelse:# the edge crosses the partition cut# ~1 ms: send a request, wait for the replyreturn rpc(machine_of(edge.other_end), edge)
def walk_cost(path):cost = 0for hop in path: # always 3 hops hereif crosses_cut(hop):cost += RPC_MS # ~1 ms network round tripelse:cost += LOCAL_US # ~1 µs pointer followreturn cost # boundary = where O(1) becomes O(network)
def edge_cut_assign(graph):for v in graph.vertices:machine[v] = pick_machine(v) # 1 vertex → 1 machinecut = []for (a, b) in graph.edges:if machine[a] != machine[b]: # endpoints splitcut.append((a, b)) # severed: a hop here is an RPCreturn machine, cut # connected graph ⇒ cut is never empty
Where this sits in Build a graph database (Neo4j / Dgraph-style)
Scene 11 of 16, in the Scale & ACID act — ACID on one box; the partition cut turns hops to RPCs.. Split a graph across two machines and any cut severs edges, so a traversal that crosses the cut turns each pointer dereference into a network round trip — the boundary is exactly where O(1) becomes O(network).
Up next. Every cut crosses some edges — the question is which cut crosses the fewest. Assign each vertex to a machine (edge-cut) and a supernode dumps all ten million of its edges on one box: a hotspot. The clever alternative splits the hub itself across machines. The supernode from earlier is back, now as the thing that breaks naive partitioning.
All 16 scenes in Build a graph database (Neo4j / Dgraph-style) · Every curriculum