Parallel Programming
Patterns
McCool et al., Chapter 3
Moreno Marzolla
Dip. di Informatica—Scienza e Ingegneria (DISI)
Università di Bologna
Copyright © 2013, 2017–2025
Moreno Marzolla, Università di Bologna, Italy
[Link]
Except where stated otherwise, this work is licensed under the Creative Commons
Attribution-ShareAlike 4.0 International License (CC BY-SA 4.0). To view a copy of this
license, visit [Link] or send a letter to Creative
Commons, 543 Howard Street, 5th Floor, San Francisco, California, 94105, USA.
Parallel Programming Patterns 2
What is a pattern?
●
A design pattern is “a general solution to a recurring
engineering problem”
●
A design pattern is not a ready-made solution to a
given problem...
●
...rather, it is a description of how a certain kind of
problem can be solved
Parallel Programming Patterns 3
Architectural patterns
●
The term “architectural
pattern” was first used by
architect Christopher
Alexander to denote
common design decision
that have been used by
architects and engineers
to realize buildings and
constructions in general Christopher Alexander,
(1936--), A Pattern Language:
Towns, Buildings, Construction
Parallel Programming Patterns 4
Example
●
Building a bridge across a river
●
You do not “invent” a new type of bridge each time
– Instead, you adapt an already existing type of bridge
Parallel Programming Patterns 5
Example
Parallel Programming Patterns 6
Example
Parallel Programming Patterns 7
Example
Parallel Programming Patterns 8
Parallel Programming Patterns
●
Embarrassingly Parallel
●
Partition
●
Master-Worker
●
Stencil
●
Reduce
●
Scan
Parallel Programming Patterns 9
Parallel programming patterns:
Embarrassingly parallel
Parallel Programming Patterns 10
Embarrassingly Parallel
●
Applies when the computation can be decomposed in independent
tasks that require little or no communication
●
Examples:
– Vector sum
– Mandelbrot set
– 3D rendering
– Brute force password cracking
– ...
Processor 0 Processor 1 Processor 2
a[]
+ + +
b[]
= = =
c[] Parallel Programming Patterns 11
Scatter-Gather
●
Scatter-Gather is a practical realization of the
Embarrassingly Parallel pattern
Scatter
P0 P1 P2 P3
Map
Gather
Parallel Programming Patterns 12
Parallel programming patterns:
Partition
Parallel Programming Patterns 13
Partition
●
The input data space (in short, domain) is split in
disjoint regions called partitions
●
Each processor operates on one partition
●
This pattern is particularly useful when the application
exhibits locality of reference
– i.e., when processors can refer to their own partition only
and need little or no communication with other processors
Parallel Programming Patterns 14
Example
●
Matrix-vector product
Ax = b
● Matrix A[][] is Core 0
partitioned into P
horizontal blocks
Core 1
● Vector x[] must be
shared/replicated on all x =
processors Core 2
●
Each processor
– operates on one block of Core 3
A[][] and on a full copy of
x[]
A[][] x[] b[]
– computes a portion of the
result b[]
Parallel Programming Patterns 15
Partition
●
Types of partition
– Regular: the domain is split into partitions of roughly the
same size and shape. E.g., matrix-vector product
– Irregular: partitions do not necessarily have the same size or
shape. E.g., heath transfer on irregular solids
●
Size of partitions (granularity)
– Fine-Grained: many small partitions
– Coarse-Grained: a few large partitions
Parallel Programming Patterns 16
1-D Partitioning
●
Block
Core 0 Core 1 Core 2 Core 3
●
Cyclic
Parallel Programming Patterns 17
2-D Block Partitioning
Block, * *, Block Block, Block
Core 0
Core 1
Core 2
Core 3
Parallel Programming Patterns 18
2-D Cyclic Partitioning
Cyclic, * *, Cyclic
Parallel Programming Patterns 19
2-D Cyclic Partitioning
Cyclic-cyclic
Parallel Programming Patterns 20
Irregular partitioning example
●
A lake surface is
approximated with a
triangular mesh
●
Colors indicate the
mapping of mesh
elements to processors
Parallel Programming Patterns 21
Index mapping
●
Usually elements have a local and global index
BLKLEN = 4
Block id 0 1 2 3
Local index 0 1 2 3 0 1 2 3 0 1 2 3 0 1 2 3
Global index 0 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
●
Assuming that all blocks have the same length
BLKLEN
Global_index = Local_index + Block_id BLKLEN
Parallel Programming Patterns 22
Computation
Fine grained vs Communication
Coarse grained partitioning
●
Fine-grained Partitioning
– Better load balancing, especially if combined
with the master-worker pattern (see later)
Time
– If granularity is too fine, the computation /
communication ratio might become too low
(communication dominates on computation)
●
Coarse-grained Partitioning
– In general improves the computation /
communication ratio
– However, it might cause load imbalance
●
The "optimal" granularity is sometimes
problem-dependent; in other cases the
Time
user must choose which granularity to
use
Parallel Programming Patterns 23
Choosing the partition size
The optimal partition size is in general system- and application-
dependent; it might be estimated by measurement
Wall-clock time
“Optimal” partition size
Partition size
Too small = higher scheduling overhead Too large = imbalanced workload
Parallel Programming Patterns 24
Example: Mandelbrot set
●
The Mandelbrot set is the
set of points c on the
complex plane s.t. the
sequence zn(c) defined as
{
z n (c)= 2
0 if n=0
z n−1 (c) + c otherwise
does not diverge when
n → +∞
Parallel Programming Patterns 25
Mandelbrot set in color
●
If |zn(c)| does not exceed 2
after maxit iterations, the
pixel is black
– the point is assumed to be
part of the Mandelbrot set
●
Otherwise, the color
depends on the number of
iterations required for |zn(c)|
to become > 2
Parallel Programming Patterns 26
Pseudocode
Embarrassingly parallel
structure: the color of each
pixel can be computed
maxit = 1000; independently from other pixels
for each point (cx, cy) {
x = y = 0;
it = 0;
while ( it < maxit AND x*x + y*y ≤ 2*2 ) {
xnew = x*x - y*y + cx;
ynew = 2*x*y + cy;
x = xnew;
y = ynew;
it = it + 1;
}
plot(cx, cy, it);
}
Source: [Link]
Parallel Programming Patterns 27
Mandelbrot set
●
A regular partitioning
can result in uneven
load distribution
– Black pixels require
maxit iterations
– Other pixels require
fewer iterations
Parallel Programming Patterns 28
Load balancing
●
Ideally, each processor should perform the same
amount of work
– If the tasks synchronize at the end of the computation, the
execution time will be that of the slower task
Task 0
Task 1 busy
Task 2
idle
Task 3
barrier synchronization
Parallel Programming Patterns 29
Load balancing HowTo
●
The workload is balanced if each processor performs
more or less the same amount of work
●
Ways to achieve load balancing:
– Use fine-grained partitioning
●
...but beware of the possible communication overhead if the tasks
need to communicate
– Use dynamic task allocation (master-worker paradigm)
●
...but beware that dynamic task allocation might incur in higher
overhead with respect to static task allocation
Parallel Programming Patterns 30
Master-worker paradigm
(process farm, work pool)
●
Apply a fine-grained partitioning
– number of task >> number of execution units
●
Each task is assigned to the first available worker
Worker
0
Worker
1
Bag of tasks of possibly
Worker
different duration
P-1
Parallel Programming Patterns 31
coarse-grained decomposition block size = 64
static task assignment static task assignment
P0
P0 P1
P2
P3
P1
P0
P1
P2 P2
P3
P0
P3
P1
block size = 64 P0
dynamic (master-worker) task assignment P1
(example) P2
P3
P0
P2
P0
P3
P2
Parallel Programming Patterns 32
P0
Example
omp-mandelbrot.c
●
Coarse-grained partitioning
– OMP_SCHEDULE="static" ./omp-mandelbrot
●
Cyclic, fine-grained partitioning (64 rows per block)
– OMP_SCHEDULE="static,64" ./omp-mandelbrot
●
Dynamic, fine-grained partitioning (64 rows per block)
– OMP_SCHEDULE="dynamic,64" ./omp-mandelbrot
●
Dynamic, fine-grained partitioning (1 row per block)
– OMP_SCHEDULE="dynamic" ./omp-mandelbrot
Parallel Programming Patterns 33
Parallel programming patterns:
Stencil
Parallel Programming Patterns 34
Stencil
●
Stencil computations involve a grid whose values are
updated according to a fixed pattern called stencil
– Example: the Gaussian smoothing of an image updates the
color of each pixel with the weighted average of the previous
colors of the 5 ´ 5 neighborhood
1 4 5 4 1
4 16 28 16 4
7 28 41 28 7
4 16 28 16 4
1 4 7 4 1
Parallel Programming Patterns 35
2D Stencils
5-point 2-axis 2D stencil 9-point 1-plane 2D stencil
(von Neumann neighborhood) 9-point 2-axis 2D stencil (Moore neighborhood)
Parallel Programming Patterns 36
3D Stencils
13-point 3-axis 3D stencil
7-point 3-axis 3D stencil
Parallel Programming Patterns 37
Stencils
●
Stencil computations usually employ two domains to
keep the current and next values
– Values are read from the current domain
– New values are written to the next domain
– current and next are exchanged at the end of each step
Current domain
Next domain
Parallel Programming Patterns 38
Ghost Cells
●
How do we handle cells on
Ghost cells
the border of the domain?
– For some applications, cells
outside the domain have
some fixed, application-
dependent value
– In other cases, we may
assume periodic boundary Domain
conditions
●
In either case, we can
extend the domain with
ghost cells, so that cells on
the border do not require
any special treatment
Parallel Programming Patterns 39
[Link]
ld-i-animate-a-plane-into-a-pipe-and-then-a-pipe-into-a-torus
Periodic boundary conditions:
How to fill ghost cells
Parallel Programming Patterns 40
……..
Periodic boundary conditions:
How to fill ghost cells
Parallel Programming Patterns 41
……..
Periodic boundary conditions:
How to fill ghost cells
Parallel Programming Patterns 42
……..
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 43
….
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 44
….
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 45
….
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 46
….
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 47
….
Periodic boundary conditions:
Another way to fill ghost cells
Parallel Programming Patterns 48
….
Parallelizing stencil computations
●
Computing the next domain from the current one has
embarrassingly parallel structure
Initialize current domain
while (!terminated) {
Init ghost cells
Compute next domain in parallel
Exchange current and next domains
}
Parallel Programming Patterns 49
Stencil computations on distributed-
memory architectures
●
Ghost cells are essential to efficiently implement
stencil computations on distributed-memory
architectures
Parallel Programming Patterns 50
Example: 2D (Block, *) partitioning with 5P stencil
Periodic boundary
P0
P1
P2
Parallel Programming Patterns 51
Example: 2D (Block, *) partitioning with 5P stencil
Periodic boundary
Parallel Programming Patterns 52
Example: 2D (Block, *) partitioning with 5P stencil
Periodic boundary
Parallel Programming Patterns 53
Example: 2D (Block, *) partitioning with 5P stencil
Periodic boundary
Parallel Programming Patterns 54
2D Stencil Example:
Game of Life
●
2D cyclic domain, each cell has two possible states
– 0 = dead
– 1 = alive
●
The state of a cell at time t + 1 depends on
– the state of that cell at time t
– the number of alive cells at time t among the 8 neighbors
●
Rules:
– Alive cell with less than 2 alive neighbors → dies
– Alive cell with two or three alive neighbors → lives
– Alive cell with more than three alive neighbors → dies
– Dead cell with three alive neighbors → lives
Parallel Programming Patterns 60
Example: Game of Life
●
See game-of-life.c
Parallel Programming Patterns 61
Parallel programming patterns:
Reduce
Parallel Programming Patterns 62
Reduce
●
A reduction is the application of an associative binary
operator (e.g., sum, product, min, max...) to the
elements of an array [x0, x1, … xn-1]
– sum-reduce( [x0, x1, … xn-1] ) = x0+ x1+ … + xn-1
– min-reduce( [x0, x1, … xn-1] ) = min { x0, x1, … xn-1}
– …
●
A reduction can be realized in O(log2 n) parallel steps
Parallel Programming Patterns 63
Example: sum-reduce
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
Parallel Programming Patterns 64
Example: sum-reduce
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
2 -2 14 8 9 6 -6 3
Parallel Programming Patterns 65
Example: sum-reduce
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
2 -2 14 8 9 6 -6 3
11 4 8 11
Parallel Programming Patterns 66
Example: sum-reduce
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
2 -2 14 8 9 6 -6 3
11 4 8 11
19 15
Parallel Programming Patterns 67
Example: sum-reduce
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
2 -2 14 8 9 6 -6 3
11 4 8 11
19 15
34
Parallel Programming Patterns 68
Example: sum-reduce
n
n2
1 3 -2 4 7 11 -8 2 1 -5 16 4 2 -5 2 1
2 -2 14 8 9 6 -6 3
11 4 8 11
int n2;
19 15 do {
n2 = (n + 1)/2;
for (int i=0; i<n2; i++) {
if (i+n2<n) x[i] += x[i+n2];
}
n = n2;
34 } while (n2 > 1)
return
Parallel Programming Patternsx[0]; 69
See reduction.c
Work efficiency
●
How many sums are computed by the parallel
reduction algorithm?
– n / 2 sums at the first level
– n / 4 sums at the second level n
– …
n/2 n/4 n/8 ...
– n / 2j sums at the j-th level
– …
– 1 sum at the (log2 n)-th level
●
Total: O(n) sums
– The tree-structured reduction algorithm is work-efficient,
which means that it performs the same amount of “work” of
the optimal serial algorithm
Parallel Programming Patterns 70
….
Parallel programming patterns:
Scan
Parallel Programming Patterns 71
Scan (Prefix Sum)
●
A scan computes all prefixes of an array [x0, x1, … xn-1]
using a given associative binary operator op (e.g.,
sum, product, min, max... )
[y0, y1, … yn - 1] = inclusive-scan( op, [x0, x1, … xn - 1] )
where
y0 = x0
y1 = x0 op x1
y2 = x0 op x1 op x2
…
yn - 1= x0 op x1 op … op xn - 1
Parallel Programming Patterns 72
Scan (Prefix Sum)
●
A scan computes all prefixes of an array [x0, x1, … xn-1]
using a given associative binary operator op (e.g.,
sum, product, min, max... )
[y0, y1, … yn - 1] = exclusive-scan( op, [x0, x1, … xn - 1] )
where
y0 = 0 this is the neutral element of
y1 = x0 the binary operator (zero for
sum, 1 for product, ...)
y2 = x0 op x1
…
yn - 1= x0 op x1 op … op xn - 2
Parallel Programming Patterns 73
Example
x[] = 1 -3 12 6 2 -3 7 -10
inclusive-scan(+, x) = 1 -2 10 16 18 15 22 12
Parallel Programming Patterns 74
Example
x[] = 1 -3 12 6 2 -3 7 -10
exclusive-scan(+, x) = 0 1 -2 10 16 18 15 22
Parallel Programming Patterns 75
Serial implementation
void inclusive_scan(int *x, int *s, int n) // n must be > 0
{
int i;
s[0] = x[0];
for (i=1; i<n; i++) {
s[i] = s[i-1] + x[i];
}
}
void exclusive_scan(int *x, int *s, int n) // n must be > 0
{
int i;
s[0] = 0;
for (i=1; i<n; i++) {
s[i] = s[i-1] + x[i-1];
}
}
Parallel Programming Patterns 76
[Link]
Exclusive scan: Up-sweep
x[0] ∑x[0..1] x[2] ∑x[0..3] x[4] ∑x[4..5] x[6] ∑x[0..7]
x[0] ∑x[0..1] x[2] ∑x[0..3] x[4] ∑x[4..5] x[6] ∑x[4..7]
x[0] ∑x[0..1] x[2] ∑x[2..3] x[4] ∑x[4..5] x[6] ∑x[6..7]
x[0] x[1] x[2] x[3] x[4] x[5] x[6] x[7]
for ( d=1; d<n/2; d *= 2 ) {
for ( k=0; k<n; k+=2*d ) {
x[k+2*d-1] = x[k+d-1] + x[k+2*d-1];
}
} O(n) additions
Parallel Programming Patterns 77
….
[Link]
Exclusive scan: Down-sweep
x[0] ∑x[0..1] x[2] ∑x[0..3] x[4] ∑x[4..5] x[6] ∑x[0..7]
zero
x[0] ∑x[0..1] x[2] ∑x[0..3] x[4] ∑x[4..5] x[6] 0
x[0] ∑x[0..1] x[2] 0 x[4] ∑x[4..5] x[6] ∑x[0..3]
x[0] 0 x[2] ∑x[0..1] x[4] ∑x[0..3] x[6] ∑x[0..5]
0 x[0] ∑x[0..1] ∑x[0..2] ∑x[0..3] ∑x[0..4] ∑x[0..5] ∑x[0..6]
x[n-1] = 0; O(n) additions
for ( ; d > 0; d >>= 1 ) {
for (k=0; k<n; k += 2*d ) {
float t = x[k+d-1];
x[k+d-1] = x[k+2*d-1];
x[k+2*d-1] = t + x[k+2*d-1];
} Parallel Programming Patterns 78
…. }
See prefix-sum.c
Example: Line of Sight
●
n peaks of heights h[0], … h[n - 1]; the distance
between consecutive peaks is one
●
Which peaks are visible from peak 0?
not
visible visible
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 79
Line of sight
Source: Guy E. Blelloch, Prefix Sums and Their Applications
Parallel Programming Patterns 80
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 81
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 82
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 83
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 84
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 85
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 86
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 87
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 88
Line of sight
h[0] h[1] h[2] h[3] h[4] h[5] h[6] h[7]
Parallel Programming Patterns 89
Serial algorithm
●
For each i = 0, … n – 1
– Let a[i] be the slope of the line connecting the peak 0 to the
peak i
– a[0] ← -∞
– a[i] ← arctan( ( h[i] – h[0] ) / i ), se i > 0
●
For each i = 0, … n – 1
– amax[0] ← -∞
– amax[i] ← max {a[0], a[1], … a[i – 1]}, se i > 0
●
For each i = 0, … n – 1
– If a[i] ≥ amax[i] then the peak i is visible
– otherwise the peak i is not visible
Parallel Programming Patterns 90
Serial algorithm
bool[0..n-1] Line-of-sight( double h[0..n-1] )
bool v[0..n-1]
double a[0..n-1], amax[0..n-1]
a[0] ← -∞
for i ← 1 to n-1 do
a[i] ← arctan( ( h[i] – h[0] ) / i )
endfor
amax[0] ← -∞
for i ← 1 to n-1 do
amax[i] ← max{ a[i-1], amax[i-1] }
endfor
for i ← 0 to n-1 do
v[i] ← ( a[i] ≥ amax[i] )
endfor
return v
Parallel Programming Patterns 91
Serial algorithm
bool[0..n-1] Line-of-sight( double h[0..n-1] )
bool v[0..n-1]
double a[0..n-1], amax[0..n-1]
a[0] ← -∞
for i ← 1 to n-1 do Embarrassingly
a[i] ← arctan( ( h[i] – h[0] ) / i ) parallel
endfor
amax[0] ← -∞
for i ← 1 to n-1 do
amax[i] ← max{ a[i-1], amax[i-1] }
endfor
for i ← 0 to n-1 do
v[i] ← ( a[i] ≥ amax[i] ) Embarrassingly
endfor parallel
return v
Parallel Programming Patterns 92
Parallel algorithm
bool[0..n-1] Parallel-line-of-sight( double h[0..n-1] )
bool v[0..n-1]
double a[0..n-1], amax[0..n-1]
a[0] ← -∞
for i ← 1 to n-1 do in parallel
a[i] ← arctan( ( h[i] – h[0] ) / i )
endfor
amax ← exclusive-scan( max, a )
for i ← 0 to n-1 do in parallel
v[i] ← ( a[i] ≥ amax[i] )
endfor
return v
Parallel Programming Patterns 93
Conclusions
●
A parallel programming patterns defines:
– a partitioning of the input data
– a communication structure among parallel tasks
●
Parallel programming patterns can help to define
efficient algorithms
– Many problems can be solved using one or more known
patterns
Parallel Programming Patterns 94