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.
| Slurm job | Layout | Wall-clock | Speedup | Efficiency | Serial f | CPU util. |
|---|---|---|---|---|---|---|
| 46094405 | 1 node × 1 core1node1core.bigTwitter.slurm | 00:11:01 | 1.00× | 100% | n/a | 98.34% |
| 46094406 | 1 node × 8 cores1node8core.bigTwitter.slurm | 00:01:41 | 6.54× | 82% | 3.2% | 87.13% |
| 46094407 | 2 nodes × 4 cores2node8core.bigTwitter.slurm | 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.
| Slurm job | Layout | Wall-clock | Speedup | Efficiency | Serial f | CPU util. |
|---|---|---|---|---|---|---|
| 45983020 | 1 node × 1 core1node1core.bigTwitter.slurm | 00:23:03 | 1.00× | 100% | n/a | 98.40% |
| 45983021 | 1 node × 8 cores1node8core.bigTwitter.slurm | 00:03:14 | 7.13× | 89% | 1.7% | 88.20% |
| 45983022 | 2 nodes × 4 cores2node8core.bigTwitter.slurm | 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”.
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
16 workers
Amdahl10.8×
Ideal16.0×
Wall-clock time (mm:ss)
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.