# Apache Spark performance benchmarks on Arm64 and x86_64 in Google Cloud

## In this learning path

- [Introduction](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/)
- [Getting started with Apache Spark on Google Axion C4A (Arm Neoverse-V2)](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/background/)
- [How to create a Google Axion C4A Arm virtual machine on GCP](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/instance/)
- [How to deploy Apache Spark on Google Axion C4A Arm virtual machines](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/spark-deployment/)
- [Apache Spark baseline testing on Google Axion C4A Arm VM](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/baseline/)
- [Apache Spark performance benchmarks on Arm64 and x86_64 in Google Cloud](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/benchmarking/)
- [Next Steps](https://learn.arm.com/learning-paths/servers-and-cloud-computing/spark-on-gcp/_next-steps/)

## How to run Apache Spark benchmarks on Arm64 in GCP
Apache Spark includes internal micro-benchmarks to evaluate the performance of core components like SQL execution, aggregation, joins, and data source reads. These benchmarks are helpful for comparing performance on x86_64 vs Arm64 platforms.

Follow the steps outlined to run Spark’s built-in SQL benchmarks using the SBT-based framework.

1. Clone the Apache Spark source code
   ```bash
   git clone https://github.com/apache/spark.git
   ```
   This clones the full Spark source code including internal test suites and the benchmarking tools.

2. Checkout the desired Spark version
   ```bash
   cd spark/ && git checkout v4.0.0
   ```
   Switch to the stable Spark 4.0.0 release, which supports the latest internal benchmarking APIs.

3. Build Spark with benchmarking profile enabled
   ```bash
   ./build/sbt -Pbenchmarks clean package
   ```
   This compiles Spark and its dependencies, enabling the benchmarks build profile for performance testing.

4. Run a built-in benchmark suite
   ```bash
   ./build/sbt -Pbenchmarks "sql/test:runMain org.apache.spark.sql.execution.benchmark.AggregateBenchmark"
   ```
   This executes the `AggregateBenchmark`, which compares performance of SQL aggregation operations (e.g., SUM, STDDEV) with and without `WholeStageCodegen`. `WholeStageCodegen` is an optimization technique used by Spark SQL to improve the performance of query execution by generating Java bytecode for entire query stages instead of interpreting them step-by-step.

## Example Apache Spark benchmark output (Arm64)
You should see output similar to:
```
[info] Running benchmark: agg w/o group
[info]   Running case: agg w/o group wholestage off
[info]   Stopped after 2 iterations, 66883 ms
[info]   Running case: agg w/o group wholestage on
[info]   Stopped after 5 iterations, 4283 ms
[info] agg w/o group:                            Best Time(ms)   Avg Time(ms)   Stdev(ms)    Rate(M/s)   Per Row(ns)   Relative
[info] ------------------------------------------------------------------------------------------------------------------------
[info] agg w/o group wholestage off                      32967          33442         672         63.6          15.7       1.0X
[info] agg w/o group wholestage on                         856            857           1       2451.2           0.4      38.5X
```

## Understanding Apache Spark benchmark metrics and results
- **Best Time (ms):** Fastest execution time observed (in milliseconds).
- **Avg Time (ms):** Average time across all iterations.
- **Stdev (ms):** Standard deviation of execution times (lower is more stable).
- **Rate (M/s):** Rows processed per second in millions.
- **Per Row (ns):** Average time taken per row (in nanoseconds).
- **Relative Speed comparison:** baseline (1.0X) is the slower version.

## Apache Spark performance benchmark results on x86_64
The following benchmark results were collected by running the same benchmark on a `c3-standard-4` (4 vCPU, 2 core, 16 GB Memory) x86_64 virtual machine in GCP, running RHEL 9.

| Benchmark Case           | Sub-Case / Config                        | Best Time (ms) | Avg Time (ms) | Stdev (ms) | Rate (M/s) | Per Row (ns) | Relative |
|--------------------------|------------------------------------------|----------------|----------------|------------|------------|---------------|----------|
| agg w/o group            | wholestage off                           | 30044          | 32090          | 2892       | 69.8       | 14.3          | 1.0X     |
| agg w/o group            | wholestage on                            | 2728           | 2739           | 7          | 768.7      | 1.3           | 11.0X    |
| stddev                   | wholestage off                           | 4097           | 4112           | 21         | 25.6       | 39.1          | 1.0X     |
| stddev                   | wholestage on                            | 948            | 954            | 4          | 110.6      | 9.0           | 4.3X     |
| kurtosis                 | wholestage off                           | 21658          | 21664          | 9          | 4.8        | 206.5         | 1.0X     |
| kurtosis                 | wholestage on                            | 1327           | 1335           | 7          | 79.0       | 12.7          | 16.3X    |
| Aggregate w keys         | codegen = F                              | 7233           | 7234           | 1          | 11.6       | 86.2          | 1.0X     |
| Aggregate w keys         | codegen = T, hashmap = F                 | 4556           | 4570           | 21         | 18.4       | 54.3          | 1.6X     |
| Aggregate w keys         | codegen = T, row-based hashmap = T       | 1201           | 1205           | 6          | 69.8       | 14.3          | 6.0X     |
| Aggregate w keys         | codegen = T, vectorized hashmap = T      | 702            | 715            | 10         | 119.6      | 8.4           | 10.3X    |
| Aggregate w keys         | codegen = F                              | 6439           | 6524           | 119        | 13.0       | 76.8          | 1.0X     |
| Aggregate w keys         | codegen = T, hashmap = F                 | 4156           | 4170           | 12         | 20.2       | 49.5          | 1.5X     |
| Aggregate w keys         | codegen = T, row-based hashmap = T       | 2113           | 2126           | 19         | 39.7       | 25.2          | 3.0X     |
| Aggregate w keys         | codegen = T, vectorized hashmap = T      | 1310           | 1322           | 8          | 64.0       | 15.6          | 4.9X     |
| Aggregate w string key   | codegen = F                              | 2265           | 2268           | 4          | 9.3        | 108.0         | 1.0X     |
| Aggregate w string key   | codegen = T, hashmap = F                 | 1926           | 1941           | 20         | 10.9       | 91.8          | 1.2X     |
| Aggregate w string key   | codegen = T, row-based hashmap = T       | 1280           | 1285           | 8          | 16.4       | 61.0          | 1.8X     |
| Aggregate w string key   | codegen = T, vectorized hashmap = T      | 1118           | 1123           | 7          | 18.8       | 53.3          | 2.0X     |
| Aggregate w decimal key   | codegen = F                              | 2139           | 2167           | 40         | 9.8        | 102.0         | 1.0X     |
| Aggregate w decimal key   | codegen = T, hashmap = F                 | 1475           | 1488           | 18         | 14.2       | 70.3          | 1.5X     |
| Aggregate w decimal key   | codegen = T, row-based hashmap = T       | 447            | 451            | 6          | 46.9       | 21.3          | 4.8X     |
| Aggregate w decimal key   | codegen = T, vectorized hashmap = T      | 270            | 275            | 5          | 77.6       | 12.9          | 7.9X     |
| Aggregate w multiple keys | codegen = F                              | 3788           | 3834           | 65         | 5.5        | 180.6         | 1.0X     |
| Aggregate w multiple keys | codegen = T, hashmap = F                 | 2412           | 2423           | 16         | 8.7        | 115.0         | 1.6X     |
| Aggregate w multiple keys | codegen = T, row-based hashmap = T       | 1890           | 1895           | 6          | 11.1       | 90.1          | 2.0X     |
| Aggregate w multiple keys | codegen = T, vectorized hashmap = T      | 1739           | 1766           | 38         | 12.1       | 82.9          | 2.2X     |
| max func bytecode size   | codegen = F                              | 315            | 338            | 24         | 2.1        | 480.7         | 1.0X     |
| max func bytecode size   | codegen = T, hugeMethodLimit = 10000     | 178            | 200            | 13         | 3.7        | 272.3         | 1.8X     |
| max func bytecode size   | codegen = T, hugeMethodLimit = 1500      | 174            | 188            | 22         | 3.8        | 264.8         | 1.8X     |
| cube                     | wholestage off                           | 1864           | 1867           | 5          | 2.8        | 355.5         | 1.0X     |
| cube                     | wholestage on                            | 1060           | 1075           | 16         | 4.9        | 202.2         | 1.8X     |
| BytesToBytesMap         | UnsafeRowhash                            | 65             | 72             | 5          | 323.6      | 3.1           | 3.1X     |
| BytesToBytesMap         | murmur3 hash                             | 69             | 69             | 0          | 304.1      | 3.3           | 3.0X     |
| BytesToBytesMap         | fast hash                                | 41             | 42             | 1          | 517.4      | 1.9           | 5.0X     |
| BytesToBytesMap         | arrayEqual                               | 142            | 142            | 0          | 148.0      | 6.8           | 1.4X     |
| BytesToBytesMap         | Java HashMap (Long)                     | 640            | 651            | 6          | 96.9       | 10.3          | 1.0X     |
| BytesToBytesMap         | Java HashMap (two ints)                 | 89             | 93             | 2          | 235.4      | 4.2           | 2.3X     |
| BytesToBytesMap         | Java HashMap (UnsafeRow)                | 544            | 546            | 2          | 38.5       | 26.0          | 0.4X     |
| BytesToBytesMap         | LongToUnsafeRowMap (opt=false)          | 352            | 355            | 1          | 59.5       | 16.8          | 0.6X     |
| BytesToBytesMap         | LongToUnsafeRowMap (opt=true)           | 74             | 75             | 1          | 284.6      | 3.5           | 2.8X     |
| BytesToBytesMap         | BytesToBytesMap (off Heap)              | 623            | 628            | 7          | 33.7       | 29.7          | 0.3X     |
| BytesToBytesMap         | BytesToBytesMap (on Heap)               | 624            | 627            | 3          | 33.6       | 29.8          | 0.3X     |
| BytesToBytesMap         | Aggregate HashMap                        | 31             | 31             | 0          | 680.7      | 1.5           | 6.6X     |

## Apache Spark performance benchmark results on Arm64
Results from the earlier run on the `c4a-standard-4` (4 vCPU, 16 GB memory) Arm64 VM in GCP (RHEL 9):

| Benchmark Case           | Sub-Case / Config                        | Best Time (ms) | Avg Time (ms) | Stdev (ms) | Rate (M/s) | Per Row (ns) | Relative |
|--------------------------|------------------------------------------|----------------|----------------|------------|------------|---------------|----------|
| agg w/o group            | wholestage off                           | 32967          | 33442          | 672        | 63.6       | 15.7          | 1.0X     |
| agg w/o group            | wholestage on                            | 856            | 857            | 1          | 2451.2     | 0.4           | 38.5X    |
| stddev                   | wholestage off                           | 3765           | 3769           | 5          | 27.8       | 35.9          | 1.0X     |
| stddev                   | wholestage on                            | 870            | 872            | 2          | 120.6      | 8.3           | 4.3X     |
| kurtosis                 | wholestage off                           | 19114          | 19155          | 58         | 5.5        | 182.3         | 1.0X     |
| kurtosis                 | wholestage on                            | 943            | 946            | 3          | 111.2      | 9.0           | 20.3X    |
| Aggregate w/ keys        | codegen = F                              | 5401           | 5509           | 154        | 15.5       | 64.4          | 1.0X     |
| Aggregate w/ keys        | codegen = T, hashmap = F                 | 3103           | 3110           | 7          | 27.0       | 37.0          | 1.7X     |
| Aggregate w/ keys        | row-based hashmap = T                    | 1004           | 1017           | 11         | 83.5       | 12.0          | 5.4X     |
| Aggregate w/ keys        | vectorized hashmap = T                    | 707            | 711            | 3          | 118.7      | 8.4           | 7.6X     |
| Aggregate w/ string key  | codegen = F                              | 1938           | 1941           | 4          | 10.8       | 92.4          | 1.0X     |
| Aggregate w/ string key  | codegen = T, hashmap = F                 | 1208           | 1208           | 0          | 17.4       | 57.6          | 1.6X     |
| Aggregate w/ string key  | row-based hashmap = T                    | 820            | 829            | 5          | 25.6       | 39.1          | 2.4X     |
| Aggregate w/ string key  | vectorized hashmap = T                    | 756            | 756            | 0          | 27.8       | 36.0          | 2.6X     |
| Aggregate w/ decimal key  | codegen = F                              | 1878           | 1886           | 11         | 11.2       | 89.6          | 1.0X     |
| Aggregate w/ decimal key  | codegen = T, hashmap = F                 | 1116           | 1116           | 0          | 18.8       | 53.2          | 1.7X     |
| Aggregate w/ decimal key  | row-based hashmap = T                    | 411            | 423            | 11         | 51.0       | 19.6          | 4.6X     |
| Aggregate w/ decimal key  | vectorized hashmap = T                    | 278            | 280            | 2          | 75.4       | 13.3          | 6.8X     |
| Aggregate w/ multiple keys | codegen = F                              | 3272           | 3277           | 8          | 6.4        | 156.0         | 1.0X     |
| Aggregate w/ multiple keys | codegen = T, hashmap = F                 | 1802           | 1804           | 3          | 11.6       | 85.9          | 1.8X     |
| Aggregate w/ multiple keys | row-based hashmap = T                    | 1461           | 1468           | 10         | 14.4       | 69.7          | 2.2X     |
| Aggregate w/ multiple keys | vectorized hashmap = T                    | 1283           | 1285           | 3          | 16.4       | 61.2          | 2.6X     |
| Max function bytecode size | codegen = F                              | 263            | 268            | 4          | 2.5        | 401.6         | 1.0X     |
| Max function bytecode size | hugeMethodLimit = 10000                   | 143            | 148            | 8          | 4.6        | 217.4         | 1.8X     |
| Max function bytecode size | hugeMethodLimit = 1500                    | 129            | 132            | 3          | 5.1        | 196.6         | 2.0X     |
| Cube                     | wholestage off                           | 1572           | 1582           | 14         | 3.3        | 299.9         | 1.0X     |
| Cube                     | wholestage on                            | 841            | 843            | 2          | 6.2        | 160.4         | 1.9X     |
| BytesToBytesMap         | UnsafeRowhash                            | 137            | 137            | 0          | 153.6      | 6.5           | 1.0X     |
| BytesToBytesMap         | murmur3 hash                             | 48             | 48             | 0          | 440.6      | 2.3           | 2.9X     |
| BytesToBytesMap         | fast hash                                | 42             | 42             | 0          | 499.2      | 2.0           | 3.3X     |
| BytesToBytesMap         | Aggregate HashMap                        | 23             | 23             | 0          | 913.0      | 1.1           | 5.9X     |

## Apache Spark performance benchmarking comparison on Arm64 and x86_64
When you compare the benchmarking results you will notice that on the Google Axion C4A Arm-based instances:
- **Whole-stage code generation significantly boosts performance**, improving execution by up to **3×** (e.g., `agg w/o group` from 2728 ms to 856 ms).
- **Aggregation with Keys**, across row-based and non-hashmap variants deliver ~1.7–5.4× speedups.
- **Arm-based Spark shows strong hash performance**, `murmur3` and `UnsafeRowhash` on Arm-based instances are ~3×–5× faster, with the aggregate hashmap ~6× faster; the `fast hash` path is roughly on par.

Overall, when whole-stage codegen and vectorized hashmap paths are used, you should see multi-fold speedups on the Google Axion C4A Arm-based instances.
