Skip to content
Tweet Cruncher

Scaling lab · Spartan, April 2023

Eight cores, 6.5 times faster

The assignment required the same job on three Slurm layouts: 1 node × 1 core, 1 node × 8 cores and 2 nodes × 4 cores. These are the measured wall-clock times for 18.7 GB of tweets, what they say about the serial share of the program, and what more cores would have bought.

1 core
11:01
28.3 MB/s single-core scan
8 cores
1:41
1 node × 8 or 2 nodes × 4
Speedup
6.54×
82% parallel efficiency
Serial fraction
3.2%
caps speedup at 31×

Measured

The benchmark jobs

Speedup and efficiency are relative to the 1-core job; the serial fraction is the Karp–Flatt metric, which inverts Amdahl's law for a single measurement. CPU utilisation is Spartan's own job statistic.

Final submission

bigTwitter.json, jobs 46094405–07

Swipe the table sideways for efficiency, serial fraction and CPU use.

Final benchmark jobs on Spartan
LayoutWall-clockSpeedupEfficiencySerial fCPU util.
1 node × 1 core
00:11:01
1.00×100%n/a98.34%
1 node × 8 cores
00:01:41
6.54×82%3.2%87.13%
2 nodes × 4 cores
00:01:41
6.54×82%3.2%87.75%

Earlier revision

Same file and layouts on 2 April 2023, before the final optimisations (Slurm logs on the Spartan-running-test branch)

Swipe the table sideways for efficiency, serial fraction and CPU use.

Earlier benchmark jobs on Spartan
LayoutWall-clockSpeedupEfficiencySerial fCPU util.
1 node × 1 core
00:23:03
1.00×100%n/a98.40%
1 node × 8 cores
00:03:14
7.13×89%1.7%88.20%
2 nodes × 4 cores
00:03:13
7.17×90%1.7%86.50%

Model

Fit Amdahl's law, then push it

Amdahl's law says that if a fraction f of a job is inherently serial, n workers can at best speed it up by 1 / (f + (1 − f) / n). Fitted to our runs, the curve explains why eight cores gave 6.5× rather than 8×.

Amdahl's law explorer

The curve starts from the measured 1-core time and assumes a fraction f of the work can never be split. Its starting value is the least-squares fit to the 8-core runs. Drag the sliders to ask “what if”.

Fitted: 3.18%
Moves the crosshair; hovering a chart does too.

Amdahl's law with f = 3.18% predicts:

Time at n = 16
1:01
Speedup at n = 16
10.84×
Efficiency
68%
Ceiling, n → ∞
31.5×

Speedup vs 1 core

Speedup against number of workers: measured Spartan runs, Amdahl curve and ideal linear speedup0×10×20×30×40×1102030workers (MPI ranks / cores)

16 workers

Amdahl10.8×

Ideal16.0×

Wall-clock time (mm:ss)

Wall-clock time against number of workers: measured Spartan runs, Amdahl curve and ideal0:003:006:009:0012:001102030workers (MPI ranks / cores)

16 workers

Amdahl1:01

Ideal0:41

  • 1 node × 1 core
  • 1 node × 8 cores
  • 2 nodes × 4 cores
  • Amdahl's law
  • Ideal (linear)

Reading the numbers

Three things the benchmarks show

Serial work sets the ceiling

Every rank starts Python, imports polars and pandas and builds the full sal.json dictionary before it reads a byte of tweets, and the task ranks merge the partial tables at the end. That 3.2% does not shrink with more cores, so the model tops out near 31×.

Faster code scaled worse

The earlier revision took 23:03 on one core but reached 7.1× on eight (serial share 1.7%). Precompiled regexes, a reworked read loop and a new gather step made the program twice as fast, which left the fixed costs a bigger slice of each run: Amdahl's law in practice.

A second node was free

2 nodes × 4 cores ran in 1:41, the same as one 8-core node. Ranks only exchange small per-author and per-city count tables at the end, never tweets, so the slower inter-node link barely registers.

How sure is 3.2%?

Not very, and the data cannot say how unsure. The benchmark is three configurations with one run each, and both multi-core layouts used 8 cores, so the fitted serial fraction is the Karp–Flatt value at a single point (n = 8): a point estimate with no repeat from which to estimate run-to-run spread, and no way to tell Amdahl's law from any other curve through that point. Slurm's whole-second timing alone moves it between 3.08% and 3.28%; the 31× ceiling is an extrapolation. The MPI lab's benchmark repeats every configuration and reports intervals instead (method).

Want to measure this yourself? The MPI lab runs the same chunk-and-gather design on your machine's cores and fits the same curve to your runs.