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.

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
PARTITIONmachine 0machine 1AliceBobCarolDaveErinFrankGraceHeidiIvanJudyThe RockThe MatrixInceptionGothamFRIENDRATEDLIVES_INFOLLOWSOne machine: every friend-hop from Alice is a local pointer dereference — microseconds.PARTITIONEdge-cut2 machinescross-hop: +1ms RPCSPINElocal 4 · global 14Same few nodes lit — but each…
What to watch for

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.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Traversal.follow_edge
one hop: pointer dereference if local, network RPC if it crosses the cut
def follow_edge(cur, edge):
if machine_of(edge.other_end) == this_machine:
# local hop: index-free adjacency, ~1 µs
return deref_pointer(edge) # pointer follow
else:
# the edge crosses the partition cut
# ~1 ms: send a request, wait for the reply
return rpc(machine_of(edge.other_end), edge)
Traversal.walk_cost
3 hops: total time = local follows + cut-crossing RPCs
def walk_cost(path):
cost = 0
for hop in path: # always 3 hops here
if crosses_cut(hop):
cost += RPC_MS # ~1 ms network round trip
else:
cost += LOCAL_US # ~1 µs pointer follow
return cost # boundary = where O(1) becomes O(network)
Partitioner.edge_cut_assign
assign each vertex to a machine — the edges that span machines ARE the cut
def edge_cut_assign(graph):
for v in graph.vertices:
machine[v] = pick_machine(v) # 1 vertex → 1 machine
cut = []
for (a, b) in graph.edges:
if machine[a] != machine[b]: # endpoints split
cut.append((a, b)) # severed: a hop here is an RPC
return 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

Built with Arqly
Every scene in Build a graph database (Neo4j / Dgraph-style) builds on the one before it.All 16 Build a graph database (Neo4j / Dgraph-style) scenes