Amazon EMR is AWS's managed service for open-source big data engines: Apache Spark above all, plus Hive, Trino, HBase, Flink and others. It does not replace those engines. It installs a tested set of versions, wires them to S3, IAM and CloudWatch, provisions and scales the machines, and gives you an API for submitting work. What you still own is the engine's behaviour: partitioning, shuffles, memory settings and data layout.

That split is the key to operating EMR well. Most EMR incidents are Spark or YARN incidents with an AWS-shaped trigger: a Spot reclaim, a core node that ran out of disk, a bootstrap script that failed. This article explains what EMR provisions, how the three deployment options differ, how a job flows through a cluster, how to use Spot and scaling safely, and what to check when things fail. Commands use release label emr-7.14.0, listed as the latest 7.x release when this was written; check the release guide for the current label and its application versions.

Advertisement

Three ways to run EMR

EMR is really three products that share release labels and engine builds but differ in who manages the compute.

OptionYou manageEMR managesGood fit
EMR on EC2Cluster shape, instance types, scaling limits, bootstrap actionsProvisioning, application install, YARN, HDFS, stepsLong-running or heavily tuned clusters, HBase, custom software on nodes
EMR ServerlessApplication limits and job parametersWorkers, capacity, scaling, patchingBatch and ad hoc Spark or Hive jobs with variable load
EMR on EKSThe EKS cluster, node groups, namespacesSpark runtime images and job submissionTeams that already run Kubernetes and want to share it

Choose Serverless by default for new batch Spark or Hive work, because it removes cluster sizing and idle cost. Choose EC2 when you need things Serverless does not offer, such as HBase, custom agents on the nodes, specific instance types or long-lived interactive clusters. Choose EKS when Kubernetes is already your platform and you want Spark jobs to share nodes, quotas and tooling with everything else.

Anatomy of an EMR on EC2 cluster

An EC2 cluster has three node roles. The primary node (still MASTER in the API) runs the coordinating services: the YARN ResourceManager, the HDFS NameNode and application servers such as the Spark history server. Core nodes run a YARN NodeManager and an HDFS DataNode, so they provide both compute and storage. Task nodes run only a NodeManager: compute with no HDFS blocks.

That difference drives almost every operational decision. Removing a task node loses running containers, which Spark retries. Removing a core node can lose HDFS blocks, and shrinking core capacity requires HDFS decommissioning, which is slow. So core nodes should be few, stable and On-Demand, and elastic capacity should be task nodes.

An EMR on EC2 cluster: who runs what, and where the data livesPrimary nodeYARN ResourceManager, HDFS NameNodeCore nodesNodeManager + HDFS DataNodeTask nodesNodeManager only, Spot-friendlySpark executorsYARN containersEMR control planesteps, scaling, healthAmazon S3 via EMRFSinput, output, logsGlue Data Catalogtable metadataschedulesread/writelookupsprovisionsOnly core nodes hold HDFS blocks. Losing a core node can lose data; losing a task node loses only work.Durable data belongs in S3, so the cluster itself can be transient.
A Spark job on EMR on EC2. YARN schedules executors on core and task nodes; EMRFS connects Spark to S3, and the Glue Data Catalog can serve as the Hive metastore.

Data normally lives in S3, accessed through EMRFS, EMR's S3 connector. HDFS on core nodes is best treated as scratch space for intermediate data and for applications that need it, such as HBase when not configured on S3. Since S3 became strongly consistent for reads after writes, the old EMRFS consistent view feature is not needed; do not enable it on new clusters. Table metadata can come from the AWS Glue Data Catalog, which lets EMR, Athena and Glue jobs share tables; see building a data lake on AWS.

Advertisement

Release labels and configuration

A release label such as emr-7.14.0 pins a set of application versions, EMR patches and the Amazon Linux base image. Everything on the cluster comes from that label, so upgrading Spark means moving to a new label, and testing a new label is the unit of upgrade work. Pin the label in infrastructure code and change it deliberately.

Configuration is passed as classifications, which EMR maps onto each application's config files. spark-defaults becomes spark-defaults.conf, yarn-site becomes yarn-site.xml, and so on. Setting values through classifications, rather than by editing files in bootstrap actions, keeps them visible in the cluster description and applied consistently to nodes added later by scaling.

Bootstrap actions are scripts that run on every node before applications start. They are the place for OS packages and agents. Keep them short and idempotent, because a failed bootstrap action fails the node, and on the initial nodes it fails the whole cluster launch.

Worked example: a nightly Spark job on a transient cluster

Suppose a job reads one day of order events from S3, about 400 GB of compressed Parquet, joins a customer dimension and writes daily aggregates, once a night. The classic pattern is a transient cluster: create, run one step, terminate. The first block is the configs.json file.

[
  {"Classification": "spark-defaults",
   "Properties": {"spark.sql.adaptive.enabled": "true",
                  "spark.dynamicAllocation.enabled": "true"}}
]


# A transient cluster: run one step, then terminate.
aws emr create-cluster \
  --name nightly-orders \
  --release-label emr-7.14.0 \
  --applications Name=Spark \
  --use-default-roles \
  --log-uri s3://acme-emr-logs/nightly/ \
  --configurations file://configs.json \
  --instance-groups \
     InstanceGroupType=MASTER,InstanceCount=1,InstanceType=m7g.xlarge \
     InstanceGroupType=CORE,InstanceCount=2,InstanceType=r7g.2xlarge \
  --steps Type=Spark,Name=orders,ActionOnFailure=TERMINATE_CLUSTER,\
Args=[--deploy-mode,cluster,s3://acme-jobs/orders_daily.py,--date,2026-09-30] \
  --auto-terminate

Walk through the lifecycle. EMR provisions the instances, runs bootstrap actions, installs Spark and YARN from the release label, then submits the step with spark-submit in cluster mode, so the driver runs in a YARN container rather than on the primary node. The step reads from S3 through EMRFS and writes output through the EMRFS S3-optimized committer, which avoids renames on S3. When the step ends, --auto-terminate shuts the cluster down. ActionOnFailure=TERMINATE_CLUSTER makes a failed step tear the cluster down too, so a broken job does not leave instances running.

Size it with simple arithmetic first. An r7g.2xlarge has 8 vCPUs and 64 GiB. With four cores per executor and some memory left for the OS and YARN overhead, each node runs about two executors with roughly 20 GiB of heap each. Two core nodes give four executors and sixteen concurrent tasks. If the job needs about 3,200 input tasks at 128 MB each, that is 200 waves, which is too slow; adding task nodes, not core nodes, is how you shorten it. With spark.dynamicAllocation.enabled on, Spark asks YARN for more executors as the backlog grows; see Spark dynamic allocation.

Logs go to the --log-uri prefix in S3, which survives termination. Persistent application UIs let you open the Spark history server for a terminated cluster from the console, which matters for transient clusters where nothing survives on the nodes; the Spark history server explains what to look for once it is open.

Instance fleets, Spot and managed scaling

Instance groups use one instance type per group. Instance fleets let each node role list several instance types with weights, and EMR fills a target capacity from whichever types are available. For Spot this matters: the more instance types and Availability Zones you allow, the less likely one reclaim wave takes out all your capacity.

{
  "InstanceFleetType": "TASK",
  "TargetOnDemandCapacity": 0,
  "TargetSpotCapacity": 32,
  "InstanceTypeConfigs": [
    {"InstanceType": "r7g.2xlarge", "WeightedCapacity": 8},
    {"InstanceType": "r6g.2xlarge", "WeightedCapacity": 8},
    {"InstanceType": "r7g.4xlarge", "WeightedCapacity": 16},
    {"InstanceType": "r6g.4xlarge", "WeightedCapacity": 16}
  ],
  "LaunchSpecifications": {
    "SpotSpecification": {
      "AllocationStrategy": "capacity-optimized",
      "TimeoutDurationMinutes": 10,
      "TimeoutAction": "SWITCH_TO_ON_DEMAND"
    }
  }
}

A cluster uses either instance groups or instance fleets, never both, so this task fleet belongs to a fleet-based cluster, not the one above. It asks for 32 units of Spot, where an 8-vCPU type counts as 8. If Spot cannot be provisioned within ten minutes at launch, it switches to On-Demand. Put Spot on task fleets, keep the primary and core fleets On-Demand, and make sure Spark can tolerate losing executors: shuffle data on a reclaimed node must be recomputed, so very long stages on Spot can lose significant work.

Managed scaling adjusts cluster size from YARN metrics within limits you set. The limits are the main tuning knob.

aws emr put-managed-scaling-policy --cluster-id j-EXAMPLE123 --managed-scaling-policy '{
  "ComputeLimits": {
    "UnitType": "Instances",
    "MinimumCapacityUnits": 3,
    "MaximumCapacityUnits": 40,
    "MaximumOnDemandCapacityUnits": 6,
    "MaximumCoreCapacityUnits": 3
  }
}'

Here the cluster can grow to 40 instances, but at most 3 of them core nodes, so all growth lands on task nodes, and at most 6 On-Demand, so the rest must be Spot. The API also accepts InstanceFleetUnits and VCPU as unit types, plus an optional ScalingStrategy of DEFAULT or ADVANCED. Managed scaling reacts to YARN demand, so it helps jobs that use dynamic allocation and does little for jobs that request a fixed number of executors.

EMR Serverless: applications, workers and limits

EMR Serverless removes the cluster. You create an application of type Spark or Hive for a release label, and submit job runs to it. Each job run gets workers, which are units of vCPU, memory and disk, and EMR adds and removes them as the job's demand changes. You pay for worker resources while they run.

aws emr-serverless create-application \
  --name orders-etl --type SPARK --release-label emr-7.14.0 \
  --architecture ARM64 \
  --initial-capacity '{
     "DRIVER":   {"workerCount": 1, "workerConfiguration": {"cpu": "2vCPU", "memory": "4GB"}},
     "EXECUTOR": {"workerCount": 4, "workerConfiguration": {"cpu": "4vCPU", "memory": "16GB"}}}' \
  --maximum-capacity '{"cpu": "200vCPU", "memory": "800GB"}' \
  --auto-stop-configuration '{"enabled": true, "idleTimeoutMinutes": 15}'

aws emr-serverless start-job-run \
  --application-id 00fexample \
  --execution-role-arn arn:aws:iam::123456789012:role/orders-etl-job \
  --job-driver '{"sparkSubmit": {
      "entryPoint": "s3://acme-jobs/orders_daily.py",
      "entryPointArguments": ["--date", "2026-09-30"],
      "sparkSubmitParameters": "--conf spark.executor.cores=4 --conf spark.executor.memory=12g"}}' \
  --configuration-overrides '{"monitoringConfiguration":
      {"s3MonitoringConfiguration": {"logUri": "s3://acme-emr-logs/serverless/"}}}'

Three settings control behaviour and cost. Initial capacity keeps pre-initialized drivers and executors warm, so jobs start in seconds instead of waiting for workers; warm workers are billed while the application is started, so size them for the common job, not the largest. Maximum capacity caps the whole application's vCPU, memory and disk across all its concurrent jobs, which makes it the main guard against a runaway job or a bad retry loop. Auto-stop stops the application after an idle period, releasing warm capacity.

Each job run assumes an execution role, which is how Serverless controls access to S3 and the Glue Data Catalog. Give each application or pipeline its own role scoped to its buckets, rather than one broad role for everything.

EMR on EKS

EMR on EKS registers a Kubernetes namespace as a virtual cluster. Jobs are submitted with aws emr-containers start-job-run, naming the virtual cluster, an execution role, a release label and a sparkSubmitJobDriver with the entry point and parameters. EMR launches the driver as a pod, and the driver launches executor pods, using EMR's Spark runtime images.

The attraction is consolidation: Spark shares nodes, autoscaling, observability and cost allocation with other Kubernetes workloads. The price is that you now own node groups, pod scheduling, image management and Kubernetes upgrades, which EMR on EC2 and Serverless handle for you. Choose it when a platform team already runs EKS well.

Security and access

  • Roles. EC2 clusters use a service role for EMR itself and an instance profile for the nodes; everything on a node inherits the instance profile, so scope it tightly. Runtime roles for steps and Serverless execution roles give per-job permissions instead.
  • Networking. Launch clusters in private subnets with VPC endpoints for S3 and other services. Do not open the primary node's UI ports to the internet; reach UIs through a bastion, Session Manager or the console's application UIs.
  • Encryption. Security configurations turn on at-rest encryption for S3 data and local disks and in-transit encryption between nodes. Define them once and reference them from every cluster.
  • Data access. For fine-grained table permissions, integrate with Lake Formation rather than granting whole buckets. Amazon S3, in depth covers bucket policies and encryption options.

Failure modes

SymptomLikely causeFix
Cluster fails in STARTING or BOOTSTRAPPINGBootstrap action exited non-zero, or no capacity for the instance typeTest scripts on a single node, add retries for downloads, allow more instance types
HDFS missing blocks after a node lossSpot or unhealthy core nodesKeep core On-Demand and small; keep durable data in S3
Core nodes marked unhealthy, jobs stallLocal disk full from shuffle or logsLarger EBS volumes, clean up intermediate data, fix skew
Executors lost in wavesSpot reclaim on task nodesDiversify instance types, shorten stages, checkpoint long pipelines
Everything fails at onceSingle primary node lostUse three primary nodes for long-running clusters, transient clusters for batch
Serverless jobs queue or fail to scaleApplication maximum capacity reachedRaise the cap or split applications by workload
Slow writes and many small filesToo many output partitionsRepartition before writing, use a table format with compaction; see Spark with Iceberg

Trade-offs and cost

EMR's strength is control over open-source engines at EC2 prices, with Spot and Graviton instances available. Its weakness is that the control comes with responsibility: you size executors, tune shuffles and debug YARN. AWS Glue runs Spark with less to configure and less to tune, which suits simpler ETL; see Glue ETL, in depth. Heavier or long-running Spark workloads, custom software on the nodes, or non-Spark engines such as HBase and Trino point to EMR.

For cost: make batch clusters transient or use Serverless, so nothing idles. Move elastic capacity to Spot task nodes with diversified fleets. And fix the job before scaling the cluster: skewed joins, tiny files and unnecessary shuffles cost more than any instance choice saves.

What to do next

  1. Pick one recurring Spark job and decide, using the table above, whether it belongs on Serverless, EC2 or EKS.
  2. Pin the release label in infrastructure code and schedule a test run on each new label before upgrading.
  3. Move all configuration into classifications and keep bootstrap actions idempotent.
  4. On EC2, keep core nodes On-Demand and few, put Spot on a diversified task fleet, and cap core growth in managed scaling.
  5. On Serverless, set maximum capacity and auto-stop on every application, and a scoped execution role per pipeline.
  6. Send logs to S3, confirm you can open the Spark history server for a terminated cluster, and write a runbook from the failure table.
Key takeaway: Amazon EMR runs open-source engines, mostly Spark, on compute that AWS provisions and scales, in three forms: EC2 clusters, Serverless applications and EKS virtual clusters. On EC2, the difference between core nodes, which hold HDFS, and task nodes, which do not, decides where Spot and scaling are safe. Keep data in S3, clusters transient or serverless, configuration in classifications and elastic capacity on diversified Spot fleets. Set capacity limits and scoped roles everywhere, and remember that most EMR problems are Spark problems that EMR only triggers.