Bulk Synchronous Parallel, or BSP, is a model of parallel computation that Leslie Valiant proposed in 1990 in "A bridging model for parallel computation". The program runs as a sequence of supersteps. In each one, every processor computes on local data, sends messages, then waits at a global barrier. Nothing sent during a superstep can be read until the next one. That single rule removes data races. It also gives a cost formula simple enough to predict runtime on paper.
BSP is everywhere once you look for it. Pregel and its descendants, such as Apache Giraph and Spark's GraphX, run graph algorithms in supersteps. Synchronous data-parallel training, where every GPU computes gradients and then joins an all-reduce, is BSP with one superstep per step. This page covers the model, its cost formula, how to measure the machine parameters, two complete algorithms with code, a worked cost calculation, the ML connection and the failure modes.
The model: supersteps and barriers
A BSP computer has three parts: p processors, each with its own memory; a network that delivers messages between any pair; and a barrier that synchronizes all processors. A superstep has three phases.
- Local computation. Each processor works only on data in its own memory, including messages delivered at the end of the previous superstep.
- Communication. Processors send messages or write into other processors' memory. These are buffered and not visible to the receiver yet.
- Barrier synchronization. Everyone waits until all processors have finished both phases and all messages are delivered. Then the next superstep begins.
The model ignores network topology and message order inside a superstep. That is intentional. Valiant's argument was that a good bridging model lets software be written once and run predictably on many machines, just as the von Neumann model does for sequential code.
The cost formula
BSP describes a machine with four numbers: p, the processor count; r, the computation speed, often normalized away so work is counted in flop units; g, the gap, the cost per word of communication when the network is under continuous load; and l, the latency of a barrier. Both g and l are expressed in units of local operations.
In each superstep, let w be the largest local work on any processor. Let h be the largest number of words any processor sends or receives. The pattern is called an h-relation. Then:
cost(superstep) = w + h * g + l
cost(program) = W + H * g + S * l
where W = sum of w over supersteps
H = sum of h over supersteps
S = number of superstepsThis formula tells you what to optimize. W rewards load balance, because the maximum counts, not the average. H rewards communication volume and its balance: one processor receiving everything sets h for the whole superstep. S rewards fewer synchronizations, because every barrier costs l whatever the work. The three terms trade against each other. Batching messages lowers S but may raise memory. Replicating data lowers H but raises W. A good BSP algorithm states its cost as a formula in n and p and is tuned for the machine's ratio of g and l to compute speed.
Measuring g and l
Values for g and l come from measurement, not data sheets. A simple probe runs a series of supersteps, each performing a full h-relation for growing h, and fits a straight line:
def probe(bsp, h_values, reps=20):
rows = []
for h in h_values:
bsp.sync()
t0 = bsp.time()
for _ in range(reps):
for k in range(h): # each processor sends h words, round robin
bsp.put((bsp.pid + 1 + k) % bsp.nprocs, k)
bsp.sync()
rows.append((h, (bsp.time() - t0) / reps))
# least squares: time = g_seconds * h + l_seconds
g_s, l_s = fit_line(rows)
flops = measure_local_flop_rate()
return g_s * flops, l_s * flops # g and l in flop unitsRun it with the real processor count and the message sizes your program uses. Small messages often have a much higher effective g than large ones, and l grows with p, often logarithmically for tree barriers. On a single GPU, grid-wide synchronization or a kernel boundary plays the role of the barrier. The same probe idea measures what that costs relative to useful work.
Algorithm 1: prefix sum in two supersteps
Prefix sum (scan) is the standard first BSP algorithm. Given n numbers split into p blocks, produce every running total. It needs two supersteps:
- Each processor scans its own block locally (n/p work) and sends its block total to every higher-numbered processor (h = p - 1).
- Each processor adds the received totals to get its offset, then adds the offset to every element of its block (n/p plus p work).
Here is a complete simulator that runs the algorithm and charges BSP costs, so you can check the formula against the code:
class BSPMachine:
def __init__(self, p, g, l):
self.p, self.g, self.l = p, g, l
self.inbox = [[] for _ in range(p)]
self.costs = []
def superstep(self, step):
outbox = [[] for _ in range(self.p)]
work, sent = [0] * self.p, [0] * self.p
for pid in range(self.p):
msgs, w = step(pid, self.inbox[pid]) # reads only last superstep's messages
work[pid] = w
for dest, payload in msgs:
outbox[dest].append((pid, payload))
sent[pid] += 1
h = max(max(sent), max(len(b) for b in outbox))
self.costs.append(max(work) + h * self.g + self.l)
self.inbox = outbox # delivered only after the barrier
def bsp_prefix_sum(data, p, g=4, l=50_000):
m = BSPMachine(p, g, l)
blocks = [data[i * len(data) // p:(i + 1) * len(data) // p] for i in range(p)]
def step1(pid, _inbox):
acc = 0
for i, x in enumerate(blocks[pid]):
acc += x
blocks[pid][i] = acc
return [(d, acc) for d in range(pid + 1, p)], len(blocks[pid])
def step2(pid, inbox):
offset = sum(total for _src, total in inbox)
blocks[pid] = [x + offset for x in blocks[pid]]
return [], len(blocks[pid]) + len(inbox)
m.superstep(step1)
m.superstep(step2)
return [x for b in blocks for x in b], m.costsNotice the communication shape. Processor p-1 receives p-1 messages and processor 0 sends p-1, so h = p - 1. For large p, a tree-structured scan uses h = 1 per superstep but needs about log2 p supersteps. That trades (p-1)g for (log2 p)(g + l). Which wins depends on whether l is large relative to p times g, which is exactly the kind of decision the cost model makes explicit.
Worked example: 100 million elements on 64 processors
Take n = 100,000,000 numbers on p = 64 processors. Assume the probe measured g = 4 and l = 50,000 flop units; these are illustrative values for a cluster. Then n/p = 1,562,500.
| Superstep | w | h * g | l | Cost |
|---|---|---|---|---|
| 1: local scan, send totals | 1,562,500 | 63 * 4 = 252 | 50,000 | 1,612,752 |
| 2: add offsets | 1,562,500 + 63 | 0 | 50,000 | 1,612,563 |
| Total | 3,225,315 |
A sequential scan costs about 100,000,000 operations, so the predicted speedup is about 31 on 64 processors. Communication is negligible. The model shows where the other factor of 2 goes: each element is processed twice, once in the local scan and once to add the offset. That doubling is inherent to a work-efficient parallel scan, so counting operations alone, no rearrangement removes it. Memory traffic is a different matter, and on real hardware scan is limited by bandwidth. A reduce-then-scan variant computes only the block sum in superstep 1, then does one scan pass from the offset in superstep 2. Each element is read twice and written once instead of written twice. To see that gain in the model, count memory words as part of w. The general lesson is to choose the unit of w to match what actually limits your machine.
Algorithm 2: PageRank in Pregel
Pregel, described by Malewicz and colleagues at Google in 2010, applies BSP to graphs with a vertex-centric API. In each superstep every active vertex runs a compute function. It reads messages sent to it in the previous superstep, updates its value, sends messages along its edges and may vote to halt. A halted vertex wakes up if it receives a message. The job ends when all vertices have halted and no messages are in flight. PageRank looks like this:
def compute(vertex, messages, superstep, num_vertices):
if superstep >= 1:
vertex.value = 0.15 / num_vertices + 0.85 * sum(messages)
if superstep < 30:
if vertex.out_edges:
share = vertex.value / len(vertex.out_edges)
for edge in vertex.out_edges:
send_message(edge.target, share)
else:
vote_to_halt()Two Pregel features map straight onto the cost formula. Combiners merge messages bound for the same vertex on the sender side, here by summing them, which cuts h. Aggregators compute global values such as total rank or convergence error, and their results are visible in the next superstep. Graph partitioning determines both W and H. A power-law graph with a few hub vertices creates a processor that receives far more than its share, so h is set by the worst hub. For hands-on Spark equivalents, see GraphX.
BSP in machine learning systems
Synchronous data-parallel training is BSP. Each worker computes gradients on its shard of the batch (w), all-reduces them (h is about 2 times the model size for a ring all-reduce, times g), and steps. The all-reduce completes for everyone, so it acts as the barrier. The same lessons apply. One slow GPU or noisy host sets w for the whole step. Gradient bucketing overlaps communication with the backward pass, hiding part of h*g behind w.
When stragglers dominate, systems relax the barrier. Stale synchronous parallel (Ho and colleagues, 2013) lets fast workers run ahead of the slowest by at most a fixed number of steps, trading consistency for less waiting. Fully asynchronous updates go further and lose BSP's determinism and simple analysis. Parallel matrix multiplication shows the same communication trade-offs in a dense setting; see parallel matmul.
Failure modes
- Stragglers. Imbalance or a slow node stretches every superstep. Over-decompose the data so work can be rebalanced, and monitor per-processor time per superstep.
- Skewed h. A hub vertex or hot key makes one processor receive far more than the rest. Use combiners, split heavy vertices, or partition by edges.
- Mismatched barriers. In libraries such as BSPlib, every processor must reach every sync. A conditional sync on one rank hangs or corrupts the run. Keep sync calls outside data-dependent branches.
- Message buffer blow-up. All messages for a superstep are held until the barrier. A step that sends a message per edge can exhaust memory. Batch the work into more supersteps.
- Too many supersteps. Algorithms that need many rounds with little work each, such as BFS on a long path, are dominated by S times l. Consider a different algorithm or fewer, larger steps.
- Fault tolerance cost. Pregel-style systems checkpoint at superstep boundaries. A failure rolls back to the last checkpoint, so choose the interval from the failure rate.
Trade-offs
BSP buys simplicity: no races within a superstep, deterministic results and a cost model you can evaluate on paper. You pay for it with barrier latency and idle time under imbalance. LogP models communication in more detail and fits fine-grained message passing better. MapReduce is close to BSP with two fixed phases per round and storage between rounds, which costs more per round but survives failures cheaply. Asynchronous models remove barriers and make reasoning harder. For algorithms that need agreement despite failures rather than barriers, see distributed consensus.
What to do next
- Write the BSP cost of your algorithm as W + H*g + S*l in terms of n and p before coding.
- Measure g and l on your real machine with a probe at your actual message sizes.
- Run the simulator above on small inputs and check that predicted and measured costs agree.
- Log per-processor work and bytes sent and received per superstep to find imbalance and skew.
- Add combiners or partition changes when one processor sets h.
- Compare a tree-style and a direct-exchange version when l and p*g are of similar size.
- For training jobs, track per-step straggler time and consider bounded staleness only if it dominates.