AWS courseLesson 11 of 12
AWS course · Lesson 11 of 12
Amazon EMR: Spark Clusters, EMR on EKS and EMR Serverless
Run Spark on Amazon EMR: cluster node types, Spark tuning, instance fleets and Spot, bootstrap actions, EMR on EKS, EMR Serverless and how to keep EMR costs down.
On this page
Amazon EMR is AWS’s managed platform for open-source big data frameworks: Apache Spark above all, plus Hive, Trino, Flink, HBase and others. Where Glue gives you serverless Spark with few knobs, EMR gives you control: instance types, Spot capacity, cluster configuration, custom images and long-running clusters. It comes in three deployment options, and choosing between them is a common design question.
Sample code
The Python examples are small calculators for decisions you make on EMR (executor sizing and fleet capacity); they use illustrative inputs. The CLI commands and configurations need an AWS account and are not executed.
EMR clusters
What it is. An EMR cluster on EC2 is a set of instances with a chosen release label (for example emr-7.13.0), which fixes the versions of Spark, Hadoop, Hive and the other applications.
How it works. Nodes have three roles:
| Node type | Runs | Notes |
|---|---|---|
| Primary (formerly master) | YARN ResourceManager, HDFS NameNode, application UIs | One, or three for high availability |
| Core | YARN NodeManager and HDFS DataNode | Store HDFS data; losing them can lose data |
| Task | YARN NodeManager only | No HDFS; safe to add and remove, ideal for Spot |
Clusters are either long-running (shared, interactive, always on) or transient (created for a job, terminated when the steps finish). For batch pipelines, transient clusters with data in S3 (not HDFS) are the usual choice: no idle cost, a clean environment per run, and failures do not affect other jobs.
aws emr create-cluster --name nightly-orders \
--release-label emr-7.13.0 \
--applications Name=Spark \
--instance-groups InstanceGroupType=MASTER,InstanceCount=1,InstanceType=m6g.xlarge \
InstanceGroupType=CORE,InstanceCount=2,InstanceType=r6g.2xlarge \
--service-role EMR_DefaultRole --ec2-attributes InstanceProfile=EMR_EC2_DefaultRole,SubnetId=subnet-0abc123 \
--log-uri s3://example-logs/emr/ \
--steps Type=Spark,Name=orders,ActionOnFailure=TERMINATE_CLUSTER,Args=[--deploy-mode,cluster,s3://example-artifacts/jobs/orders.py] \
--auto-terminate
EMR uses two IAM roles: the service role that EMR uses to manage EC2, and the EC2 instance profile that your Spark code runs as when it reads S3 or calls other services. Permission errors in a job are almost always about the instance profile (or the runtime role, if you use them).
Pitfalls.
- Storing important data in HDFS on a transient cluster; it disappears at termination. Use S3 as the durable store.
- Long-running clusters that nobody turns off. Set an idle auto-termination timeout.
- Putting core nodes on Spot: a reclaimed core node takes HDFS blocks and shuffle data with it.
In interviews. Explain the three node types, why S3 rather than HDFS is the system of record, and why transient clusters per pipeline are common.
Spark on EMR
What it is. EMR’s Spark is an Amazon-optimised build of Apache Spark running on YARN, with the EMRFS connector reading and writing S3 through s3:// paths and the Glue Data Catalog available as the Hive metastore.
How it works. Submit work as an EMR step (spark-submit run by EMR), through spark-submit on the primary node, from a notebook in EMR Studio, or from an orchestrator (Step Functions has a direct EMR integration; Airflow has EMR operators). EMR 7.x releases ship Spark 3.5; the emr-spark-8.x “AWS runtime for Apache Spark” releases bring Spark 4.0 to EMR on EC2, EMR on EKS and EMR Serverless.
Executor sizing is the classic tuning task. A common rule: leave one core per node for the OS and daemons, use about four or five cores per executor, and split the memory YARN offers between executors, keeping about 10% for overhead. The inputs below are illustrative; read the real YARN memory per node from the cluster:
def executors_per_node(node_vcpus, yarn_memory_gib, cores_per_executor=4, overhead_fraction=0.10):
"""Classic Spark-on-YARN sizing: leave one core for the OS and daemons, then pack executors."""
usable_cores = node_vcpus - 1
n = usable_cores // cores_per_executor
mem_per_executor = yarn_memory_gib / n
heap = mem_per_executor / (1 + overhead_fraction)
return n, round(heap, 1), round(mem_per_executor - heap, 1)
# yarn_memory_gib is what YARN offers on the node (less than the instance's RAM); read it from the cluster.
for vcpus, yarn_mem in [(16, 56), (32, 120)]:
n, heap, overhead = executors_per_node(vcpus, yarn_mem)
print(f"{vcpus} vCPU node, {yarn_mem} GiB for YARN -> {n} executors x 4 cores, "
f"--executor-memory {heap}g + ~{overhead}g overhead")
16 vCPU node, 56 GiB for YARN -> 3 executors x 4 cores, --executor-memory 17.0g + ~1.7g overhead
32 vCPU node, 120 GiB for YARN -> 7 executors x 4 cores, --executor-memory 15.6g + ~1.6g overhead
EMR can also do this for you: the maximizeResourceAllocation setting in the spark configuration classification sizes executors to the instance type, and dynamic allocation (on by default) adds and removes executors with load. Configuration is passed as JSON classifications:
[
{
"Classification": "spark",
"Properties": { "maximizeResourceAllocation": "false" }
},
{
"Classification": "spark-defaults",
"Properties": {
"spark.sql.adaptive.enabled": "true",
"spark.dynamicAllocation.enabled": "true",
"spark.executor.cores": "4",
"spark.executor.memory": "17g"
}
},
{
"Classification": "spark-hive-site",
"Properties": {
"hive.metastore.client.factory.class": "com.amazonaws.glue.catalog.metastore.AWSGlueDataCatalogHiveClientFactory"
}
}
]
Pitfalls.
- One huge executor per node (long garbage collection pauses) or one-core executors (no task parallelism within an executor and more overhead).
- Ignoring the Spark UI and the history server, which EMR keeps for finished applications; most slow jobs are skew or shuffle problems, not too few nodes.
- Spark 4 changes defaults (for example ANSI SQL mode); test before moving to an
emr-spark-8.xrelease.
In interviews. Walk through executor sizing from a node’s cores and memory, mention dynamic allocation and AQE, and say that the Glue Data Catalog can be EMR’s metastore so tables are shared with Athena.
Instance fleets and Spot
What it is. Instance groups use one instance type per node group. Instance fleets let each node type draw from a list of instance types and purchase options (On-Demand and Spot) to meet a target capacity. Spot Instances use spare EC2 capacity at a large discount, but AWS can reclaim them with a two-minute warning.
How it works. In a fleet you set a target in capacity units (by default, units are vCPUs, or you assign a weighted capacity per type), list several types, and pick an allocation strategy, such as capacity-optimised for Spot (fewer interruptions) or lowest price or prioritised for On-Demand. Offering several similar types makes it far more likely that EMR finds Spot capacity. This model fills a target of 96 units from a priority list when each type has limited Spot availability:
def fill_fleet(target_units, options):
"""Instance fleets: meet a target in capacity units from several instance types.
options: (instance_type, weighted_capacity, available_spot_instances), in priority order."""
plan, remaining = [], target_units
for itype, weight, available in options:
if remaining <= 0:
break
count = min(available, -(-remaining // weight)) # ceiling division
if count:
plan.append((itype, count, count * weight))
remaining -= count * weight
return plan, max(remaining, 0)
options = [("r6g.4xlarge", 16, 3), ("r5.4xlarge", 16, 2), ("r6g.8xlarge", 32, 4)]
plan, short = fill_fleet(96, options)
for itype, count, units in plan:
print(f"{itype:<12} x{count} = {units} units")
print("units still missing:", short)
r6g.4xlarge x3 = 48 units
r5.4xlarge x2 = 32 units
r6g.8xlarge x1 = 32 units
units still missing: 0
A safe design: primary and core nodes On-Demand, task nodes mostly Spot. Managed scaling then grows and shrinks the cluster between limits you set, based on YARN metrics, and you can cap how much of it is On-Demand.
[
{
"InstanceFleetType": "TASK",
"TargetSpotCapacity": 96,
"InstanceTypeConfigs": [
{ "InstanceType": "r6g.4xlarge", "WeightedCapacity": 16 },
{ "InstanceType": "r5.4xlarge", "WeightedCapacity": 16 },
{ "InstanceType": "r6g.8xlarge", "WeightedCapacity": 32 }
],
"LaunchSpecifications": {
"SpotSpecification": {
"TimeoutDurationMinutes": 20,
"TimeoutAction": "SWITCH_TO_ON_DEMAND",
"AllocationStrategy": "capacity-optimized"
}
}
}
]
Pitfalls.
- Spot for core nodes or for the primary node of a long job: an interruption can fail the whole job or lose HDFS data.
- A single Spot instance type: when that pool runs dry, the cluster cannot reach capacity.
- Long stages with huge shuffles on Spot lose a lot of work when nodes go. Spark recomputes lost partitions, which is slow but correct.
In interviews. Explain where Spot is safe (task nodes, stateless compute), why diversify instance types, and the two-minute interruption notice. Mention managed scaling with On-Demand limits.
Bootstrap actions
What it is. A bootstrap action is a script EMR runs on every node when it launches, before the applications are installed and before the node starts processing. Use it to install OS packages, Python libraries or agents, or to change system settings.
How it works. Store the script in S3 and reference it at cluster creation; you can restrict actions to the primary node by checking the instance metadata inside the script. Scripts run as the hadoop user and can use sudo.
#!/bin/bash
# s3://example-artifacts/bootstrap/install-libs.sh
set -euo pipefail
sudo python3 -m pip install --quiet "pyarrow==21.0.0"
aws emr create-cluster ... \
--bootstrap-actions Path=s3://example-artifacts/bootstrap/install-libs.sh,Name=install-libs
Pitfalls.
- A failing bootstrap action fails the node, and on startup the whole cluster. Pin versions and keep scripts small.
- Downloading packages from the internet on every launch: slow and fragile in private subnets. Use a custom AMI, an internal package mirror, or (for Python) a packaged virtual environment passed with the job.
- Changing Spark settings in bootstrap actions instead of configuration classifications.
In interviews. Know what bootstrap actions are for, that they run on every node before applications, and the alternatives (custom AMIs, configuration classifications, packaged environments).
EMR on EKS
What it is. EMR on EKS runs EMR’s Spark runtime as pods on an Amazon EKS (Kubernetes) cluster you already operate, instead of on dedicated EMR clusters.
How it works. You register a Kubernetes namespace as an EMR virtual cluster, create an IAM job execution role (mapped to pods through IAM roles for service accounts), and submit jobs with a release label. The driver and executors start as pods, scale with the Kubernetes cluster (for example with Karpenter or Cluster Autoscaler), and disappear when the job ends.
aws emr-containers start-job-run \
--virtual-cluster-id abcd1234efgh5678ijkl \
--name orders-daily \
--execution-role-arn arn:aws:iam::111122223333:role/emr-eks-job \
--release-label emr-7.13.0-latest \
--job-driver '{"sparkSubmitJobDriver": {"entryPoint": "s3://example-artifacts/jobs/orders.py", "sparkSubmitParameters": "--conf spark.executor.instances=10 --conf spark.executor.cores=4 --conf spark.executor.memory=16g"}}' \
--configuration-overrides '{"monitoringConfiguration": {"s3MonitoringConfiguration": {"logUri": "s3://example-logs/emr-eks/"}}}'
When it fits. The organisation already runs EKS and wants Spark to share that capacity, tooling, and security model; teams want different EMR versions side by side on one cluster; or jobs need custom container images.
Pitfalls.
- You operate Kubernetes: node groups, autoscaling, networking and upgrades are your job.
- Shuffle data lives on pod storage; losing nodes mid-job costs recomputation, as on EC2.
In interviews. Present EMR on EKS as “EMR’s Spark runtime on your Kubernetes”, worthwhile when Kubernetes is already the platform standard.
EMR Serverless
What it is. EMR Serverless runs Spark and Hive jobs without any cluster to configure. You create an application with a release label, submit job runs, and EMR Serverless provisions workers for each job and releases them afterwards.
How it works.
- Billing is for the vCPU, memory and storage that workers use, per second, while they run.
- Pre-initialised capacity keeps some drivers and workers warm so jobs start in seconds, at a cost while idle.
- Maximum capacity limits for the application cap spending.
- The application starts automatically on job submission and stops after a configurable idle timeout.
aws emr-serverless create-application --name orders-spark --type SPARK \
--release-label emr-7.13.0 \
--maximum-capacity '{"cpu": "200 vCPU", "memory": "800 GB"}' \
--auto-stop-configuration '{"enabled": true, "idleTimeoutMinutes": 15}'
aws emr-serverless start-job-run --application-id 00abcdef12345678 \
--execution-role-arn arn:aws:iam::111122223333:role/emr-serverless-job \
--job-driver '{"sparkSubmit": {"entryPoint": "s3://example-artifacts/jobs/orders.py", "sparkSubmitParameters": "--conf spark.executor.cores=4 --conf spark.executor.memory=16g"}}'
| Option | You manage | Best for |
|---|---|---|
| EMR on EC2 | Clusters, instance types, Spot, bootstrap, scaling | Long-running or heavily tuned workloads, many frameworks, maximum control |
| EMR on EKS | Kubernetes cluster | Organisations standardised on Kubernetes |
| EMR Serverless | Almost nothing (application limits) | Intermittent Spark or Hive batch jobs, teams without platform engineers |
| Glue for Spark | Almost nothing (job settings) | Spark ETL with Data Catalog, bookmarks, crawlers and data quality built in |
Pitfalls.
- Assuming serverless means cheapest. For steady, all-day workloads, a well-run EC2 cluster with Spot can cost less.
- No maximum capacity, so a runaway job scales as far as the account quotas allow.
- Network access to private resources needs the application to be configured for your VPC.
In interviews. Compare the four Spark options in the table by control versus effort, and pick by workload shape: intermittent jobs point to Serverless or Glue, steady heavy use points to EC2 clusters with Spot.
Cost optimisation on EMR
What it is. EMR on EC2 costs are the EC2 instances (and EBS) plus a per-second EMR charge per instance; EMR on EKS adds a charge for the vCPU and memory jobs use; EMR Serverless charges for worker vCPU, memory and storage. The design levers are the same across all three.
The levers.
- Transient clusters that terminate when the work is done, or idle auto-termination for interactive ones.
- Spot for task nodes with diversified instance fleets; On-Demand for primary and core.
- Managed scaling with sensible minimums and maximums.
- Graviton instances (the
gfamilies) where your libraries support arm64, often better price-performance. - Right-size executors from Spark UI metrics instead of over-provisioning memory.
- Read less: Parquet, partition pruning, and compacted files, exactly as for Athena and Glue.
- Newer releases: EMR’s optimised Spark runtime improves with releases; test upgrades for speed-ups.
- Tag clusters and applications for cost allocation, and alarm on clusters running longer than expected.
- Savings Plans or Reserved Instances for steady baseline capacity that runs all the time.
Pitfalls.
- Scaling out to fix a skewed stage; one task still holds everything. Fix the skew.
- Forgetting EBS volumes, cross-AZ data transfer and NAT gateway traffic, which add up on large clusters.
In interviews. Give the levers in priority order with the reason for each, and show you know which options fit which workload shape.
Practice questions
Why are task nodes good candidates for Spot Instances but core nodes are not?
Task nodes run only YARN containers and hold no HDFS data, so losing one costs only recomputation of its tasks. Core nodes store HDFS blocks (and often shuffle data); losing them can lose data or fail the job. Keep the primary and core nodes On-Demand and put elastic task capacity on Spot.
How would you size executors on a node with 16 vCPUs and 56 GiB available to YARN?
Leave one core for the OS and daemons, giving 15 usable cores; with four cores per executor that is three executors per node. Split the 56 GiB three ways (about 18.7 GiB each) and keep roughly 10% for overhead, so about 17 GiB of executor memory and about 1.7 GiB overhead. Then check the Spark UI for spills and GC time and adjust.
When would you pick EMR Serverless over Glue, and EMR on EC2 over EMR Serverless?
EMR Serverless over Glue when you want EMR’s runtime and release choices, Hive as well as Spark, or plain Spark without Glue-specific libraries, while still avoiding clusters. EMR on EC2 over Serverless for steady, heavy workloads where Spot and Reserved capacity make clusters cheaper, for frameworks Serverless does not run (HBase, Trino, Flink on EMR), or when you need bootstrap actions, custom AMIs or detailed cluster tuning.
A bootstrap action that installs Python packages makes cluster start-up fail intermittently. What do you do?
The script probably depends on internet downloads that time out or on unpinned versions that changed. Pin versions, mirror packages internally or bake them into a custom AMI, or ship a packaged virtual environment with the job (for example with --archives). Make the script fail fast with clear logging, and check the bootstrap logs in the cluster’s S3 log location.
List the main ways to reduce the cost of a nightly EMR pipeline.
Transient clusters that auto-terminate, Spot task fleets with diversified instance types, managed scaling, Graviton instances, right-sized executors, Parquet with partition pruning and compaction, tagging for cost allocation, and alarms on unusually long runs. For intermittent jobs, compare with EMR Serverless.
Key takeaways
- EMR on EC2 clusters have primary, core (HDFS) and task nodes; keep durable data in S3 and prefer transient clusters for batch.
- Size executors from cores and YARN memory, and use dynamic allocation, AQE and the Spark UI to tune.
- Instance fleets with several instance types make Spot capacity reliable; keep Spot to task nodes.
- Bootstrap actions run on every node before applications; pin and minimise them.
- EMR on EKS runs EMR Spark on your Kubernetes; EMR Serverless removes clusters entirely.
- EMR 7.x runs Spark 3.5, and the emr-spark-8.x runtime brings Spark 4.0 to all three options.
Progress is saved in this browser only. No account needed.