Lesson 302 · AWS Learning Path

AWS 302: Amazon EMR

· Published · 17 min read

Labelled process diagram for AWS 302: Job, framework, and data to Selected EMR deployment model to Managed compute execution to Output, logs, performance, and termination evidence, with decision, proof and rejection...

Why this lesson matters

Amazon EMR is not one cluster product. It is a family of ways to run open-source analytics frameworks such as Apache Spark. EMR on EC2 gives deep cluster and instance control. EMR Serverless provisions workers for applications and jobs without exposing a cluster. EMR on EKS runs EMR-managed jobs against an existing Amazon EKS platform. Selecting only by the word “serverless,” an existing Kubernetes mandate, or the lowest instance price usually ignores startup, isolation, release compatibility, storage, network paths, operator skill, interruption behavior, and idle cost.

A Spark program also behaves differently from a normal Linux process. A driver builds and coordinates work; executors process partitions; shuffles move intermediate data; skew can leave one task running long after the rest; retries can duplicate unsafe writes; and memory failures can appear far from the line that created the pressure. EMR manages infrastructure, but the architect still owns these distributed-system consequences.

This workshop teaches the common execution model first, then compares all three deployment choices. You will build a decision dossier for three workloads, size capacity, map identity and network paths, design Spot and retry safety, interpret supplied failure evidence, model cost, and write termination/cleanup proof. The required path creates no cloud resources. An optional approved path runs one tiny EMR Serverless Spark job.

What you will be able to do

By the end, you can:

  • explain Spark driver, executor, stage, task, partition, shuffle, checkpoint, and output-commit behavior;
  • distinguish EMR release labels from Spark and other application versions;
  • compare EMR on EC2, EMR Serverless, and EMR on EKS from requirements;
  • design primary, core, and task capacity without treating all nodes as interchangeable;
  • choose instance groups or instance fleets, On-Demand or Spot, and manual or managed scaling;
  • bound EMR Serverless initial and maximum capacity, auto-start, auto-stop, and job concurrency;
  • explain EKS virtual clusters, namespaces, execution roles, scheduling, and platform ownership;
  • map service, runtime, instance, Kubernetes, S3, Glue, Lake Formation, KMS, and network permissions;
  • keep durable data and recovery evidence outside ephemeral compute;
  • diagnose provisioning, bootstrap, dependency, memory, skew, shuffle, network, and output failures;
  • model complete cost rather than comparing one vCPU or EC2 rate; and
  • defend a deployment decision, operational runbook, and cleanup result.

Before you start

  • Complete AWS300 and AWS301 so that S3 layout, Glue metadata, Athena, IAM, and Lake Formation are familiar.
  • Use the T0 path unless a learner-owned account, budget, runtime role, logs bucket, and cleanup owner are approved.
  • The example Region is ap-south-1. Verify release and feature availability in the chosen Region.
  • Never launch a cluster merely to inspect the console. EMR, EC2, EBS, NAT, logs, and public IPv4 can all incur charges.
  • Do not use aws emr create-default-roles casually in an established account. It creates IAM resources with account-wide implications.
  • Redact account numbers, role ARNs, bucket names, private addresses, customer data, and signed URLs.
  • Do not store the only copy of durable output in HDFS, container disk, or node-local storage.

Create a local evidence directory and record the current caller:

mkdir -p "$HOME/aws302-evidence"
cd "$HOME/aws302-evidence"

export AWS_DEFAULT_REGION="ap-south-1"
aws sts get-caller-identity --query Arn --output text
aws configure list

1. Learn the execution model before choosing a product

submitted Spark application
          |
          v
driver: builds plan, requests resources, coordinates stages and retries
          |
          +------------+------------+
          v            v            v
      executor      executor      executor
      tasks/data    tasks/data    tasks/data
          \            |            /
           +------ shuffle --------+
                       |
                       v
        durable output, logs, metrics, and checkpoints

A job contains actions and transformations. Spark builds a directed acyclic graph, divides it into stages at shuffle boundaries, and creates tasks for data partitions. More partitions can improve parallelism until scheduling, object count, and shuffle overhead dominate. Too few partitions underuse capacity. A skewed key can make one partition enormous.

The driver is a control dependency. If it is undersized, loses network access, or terminates, adding executors does not solve the problem. Executors need CPU, heap, off-heap and overhead memory, local disk, network, and access to source/output systems. “Executor lost” is a symptom, not a diagnosis.

EvidenceWhat it helps answer
Spark event log and UIstage duration, skew, spill, failed tasks, executors
Driver logplanning, dependency, permission, serialization, final failure
Executor logtask exception, out-of-memory, disk, fetch, native library
YARN or Kubernetes eventsscheduling, eviction, container/node behavior
EMR job/step stateplatform lifecycle and failure reason
CloudWatch metricscapacity, health, pending work, storage, network
S3 inventory and checksumwhether durable output is complete

Retries are safe only when reads and writes are replay-safe. Use unique run IDs, temporary output, atomic or transactional commit patterns where supported, deterministic keys, and reconciliation. A successful platform state does not prove exactly-once business output.

2. Treat an EMR release as a software bill of materials

An EMR release label selects tested versions of Spark, Hadoop, Hive, Hudi, Iceberg, operating-system components, connectors, and EMR patches. “Spark 3” is not precise enough. Code, JARs, Python packages, table formats, serializers, JVM behavior, and configuration keys can change across releases.

As of this review, the current documentation lists EMR 7.13.0 as the latest 7.x release, while older maintained branches have different latest labels. Region rollout can lag. Never derive a production label from this sentence alone. Query the target Region and consult release notes at implementation time.

aws emr list-release-labels --max-results 20 --output json \
  | tee emr-release-labels.json

aws emr-serverless list-release-labels --max-results 20 --output json \
  | tee emr-serverless-release-labels.json

Pin the release label in infrastructure code. Maintain an upgrade test set with representative inputs, schema changes, dependencies, output checksums, performance baseline, security configuration, and rollback. Standard, extended, bridge, end-of-support, and end-of-life windows are lifecycle constraints, not documentation trivia.

3. Compare the three deployment models

DimensionEMR on EC2EMR ServerlessEMR on EKS
Primary abstractionClusterApplication and job runVirtual cluster mapped to EKS namespace
Infrastructure controlHighestLowestShared with EKS platform
Framework scopeBroad EMR applicationsSupported Serverless application typesSupported EMR container releases
StartupCluster provisioning/bootstrapWorker provisioning; warm capacity optionalPod/image/scheduler and EKS capacity path
ScalingManual, custom, or EMR managedAutomatic workers within maximumSpark dynamic allocation plus Kubernetes capacity
IsolationCluster/account/network boundariesApplication, runtime role, VPC settingsNamespace, RBAC, IAM, nodes, runtime classes
CustomizationBootstrap actions, AMI and configuration optionsRuntime configuration and job dependenciesContainer image and EKS platform controls
Idle costCluster instances remain until terminationNo active workers unless pre-initialized; other resources remainExisting EKS and node capacity can remain
Operator burdenEC2, YARN, OS, fleet, patch/release operationsJob/application operationsEMR plus Kubernetes platform operations
Best fitDeep control, broad apps, long-lived or specialized clustersIntermittent Spark/Hive with low cluster burdenExisting mature EKS platform and Kubernetes integration

Do not select EMR on EKS simply because the organization has EKS. Confirm platform tenancy, namespace isolation, scheduling capacity, image governance, upgrades, observability, and which team owns an overnight failed job.

4. Design EMR on EC2 correctly

An EMR on EC2 cluster commonly uses:

  • primary nodes for cluster coordination;
  • core nodes for compute and HDFS data; and
  • task nodes for compute without HDFS data-node responsibility.

Task nodes are usually the safer place for interruptible capacity. Losing core capacity can affect HDFS and shuffle state. Durable source, output, checkpoints, logs, scripts, and bootstrap artifacts should live in S3 or another resilient service.

Instance groups use one type and market option per group. Instance fleets allow multiple instance types and separate On-Demand and Spot target capacity. Fleets improve capacity flexibility but require weighted-capacity reasoning, subnet IP capacity, diversification, and clear allocation strategies.

DecisionSafer starting direction
Primary capacityOn-Demand; consider multi-primary only for justified long-running availability
Core capacityStable On-Demand baseline sized for required state and shuffle
Task capacityDiversified Spot where replay and deadlines tolerate interruption
DataS3 for durable data; HDFS/local disk for temporary processing
ScalingBounded managed scaling after workload tests
TerminationAuto-termination for transient clusters plus alarm and owner
ProtectionTermination protection for critical long-lived clusters, not as backup

Termination protection prevents ordinary EMR termination until disabled, but it does not preserve node-local data against every failure or manual EC2 action. It also does not make HDFS a backup.

Managed scaling observes workload and changes capacity within configured limits. Minimum, maximum, core limits, On-Demand limits, node labels, decommissioning behavior, subnet IPs, EC2 quotas, and Spot capacity all affect outcomes. Scaling cannot repair a serial algorithm, skewed key, blocked external database, or undersized driver.

5. Design EMR Serverless as bounded applications and jobs

An EMR Serverless application selects a release and application type. A job uses a runtime role and creates driver and worker capacity. Auto-start can start a stopped application when a job arrives. Auto-stop can release capacity after idle time. Maximum capacity bounds aggregate CPU, memory, and disk.

Pre-initialized capacity is a warm pool for faster startup. It is billed while the application is started even when workers are idle. Use it only when latency value exceeds idle cost, and retain auto-stop unless there is a documented reason not to.

submitter IAM identity
       |
       v
EMR Serverless application and job API
       |
       +--> job runtime role --> S3, Glue, Lake Formation, KMS
       |
       +--> optional VPC subnets/security groups --> private dependencies
       |
       v
driver and workers scale up to application maximum
       |
       v
S3/CloudWatch logs, durable output, job state, cost evidence

Maximum capacity is a financial and dependency safety control, but it must align with worker sizes. A 3-vCPU remainder cannot schedule a 4-vCPU worker. Bound concurrency and downstream connections as well as compute. Fifty concurrent Spark jobs can overwhelm a database even if EMR has capacity.

Use the job runtime role for data access, not the human submitter's credentials. Scope it to exact input, output, Glue/Lake Formation resources, KMS keys, logs, and secrets. VPC access is needed for private resources; then subnet address capacity, security groups, routes, DNS, NAT or endpoints, and source-system allowlists become part of the job.

6. Understand EMR on EKS ownership boundaries

An EMR on EKS virtual cluster registers an EKS namespace with EMR. It does not create a separate EKS control plane. A job run uses an EMR release, execution role, Spark configuration, and Kubernetes resources scheduled into that environment.

job submitter
    |
    v
EMR Containers API and virtual cluster
    |
    +--> IAM job execution role and role trust
    +--> EKS namespace, RBAC, admission policy, quota
    +--> scheduler, nodes/Fargate, CNI IPs, storage, DNS
    +--> image registry and optional custom image
    |
    v
Spark driver/executor pods --> data, logs, and outputs

The EMR team and EKS platform team share failure ownership. EMR can accept a job while Kubernetes leaves pods pending because of quotas, taints, affinity, unavailable node types, CNI address exhaustion, admission rejection, image pull failure, or autoscaler delay.

The job execution role controls AWS-resource access for job workloads. Kubernetes RBAC controls operations in the cluster. The submitter needs permission to use the virtual cluster and pass the approved role. Keep these identities separate and constrain confused-deputy paths. Custom images need provenance, vulnerability management, compatible base requirements, signing policy, and rebuild ownership.

7. Map storage, catalog, and table-format behavior

EMR workloads commonly read S3 through EMRFS or framework connectors and use the Glue Data Catalog. Lake Formation can add governed access for supported engines and modes. KMS keys must trust the effective runtime path.

Choose file and table format from update semantics:

NeedDirection
Immutable batch analyticsPartitioned Parquet or ORC
ACID table updates and snapshotsEvaluate Iceberg, Hudi, or Delta support for exact release/engine
Temporary shuffleLocal/EBS/container ephemeral storage
Durable checkpointS3 path separated by application and run
Shared metadataGlue Data Catalog with explicit schema ownership

Open table formats add metadata files, snapshots, compaction, vacuum/expiration, concurrency rules, and engine compatibility. “Stored in S3” does not remove database-like operations. Never let two engines write concurrently until their interoperability is proven.

Small files increase listing, planning, and task overhead. Giant files reduce parallelism. Partition by common selective predicates with controlled cardinality. Repartitioning, coalescing, compaction, and adaptive query execution should be tested with real distributions.

8. Build identity, encryption, and network controls

EMR on EC2 can involve the EMR service role, EC2 instance profile, per-step runtime roles where supported, auto-scaling role, and human/operator identity. Serverless and EKS use job runtime/execution roles. Each role has a different trust path and should not be collapsed into one administrator role.

Security configurations can govern in-transit and at-rest encryption for EMR on EC2. S3 bucket encryption, EBS encryption, local disk, HDFS, Kerberos, TLS certificates, KMS key policy, Secrets Manager, and application-level authentication remain distinct. Never put database passwords in bootstrap scripts, command arguments, or logs.

For private subnets, draw every dependency:

compute subnet
  -> S3 data/logs/scripts
  -> Glue and Lake Formation APIs
  -> KMS, STS, CloudWatch, Secrets Manager
  -> package/image repositories
  -> source/target database
  -> EMR service path required by deployment model

Resolve each path through a VPC endpoint, NAT, transit path, or approved public route. Include DNS, endpoint policies, security groups, NACLs, proxy rules, MTU, route symmetry, subnet IP capacity, and cross-AZ or transfer cost. A private subnet without an egress design is isolation and also a likely bootstrap failure.

9. Plan reliability and safe retries

Classify state before deciding recovery:

  • immutable source data can be reread;
  • shuffle and executor cache can usually be recomputed;
  • driver state may require whole-job retry;
  • streaming checkpoints define recovery position;
  • external side effects can duplicate;
  • partial output may need quarantine or atomic commit; and
  • HDFS-only data can disappear with cluster loss.

Spot interruption is not inherently unsafe if capacity is diversified, critical roles stay stable, tasks are replayable, and the deadline tolerates replacement. Use interruption-aware allocation and decommissioning supported by the chosen release. Test forced executor loss and capacity shortage.

For streaming, document checkpoint location, trigger interval, source offset semantics, watermark and late-data policy, state-store growth, sink idempotency, schema evolution, and recovery from a corrupted or incompatible checkpoint. A restart that processes data again can be correct only if the sink handles it.

For batch, use run manifests with input snapshot, code/version, release, parameters, output prefix, row counts, checksums, quality results, and terminal state. Publish final output only after validation.

10. Observe performance before resizing

Start with the critical path:

  1. prove input volume, object count, format, compression, and partition pruning;
  2. inspect stage timeline and identify the longest stage;
  3. compare task durations and input/shuffle sizes for skew;
  4. inspect driver and executor CPU, memory, garbage collection, spill, disk, and network;
  5. identify pending capacity or scheduler delay;
  6. check external-service latency and throttling;
  7. change one variable and rerun the same workload; and
  8. compare duration, correctness, resource time, and cost.
SymptomLikely directions to test
One task runs far longerkey skew, unsplittable input, partition imbalance
Executors repeatedly lostmemory overhead, node loss, Spot, disk, network
Driver out of memorylarge collect/broadcast, metadata, too many tasks
Heavy disk spillinsufficient execution memory or wide shuffle
Low CPU and high durationI/O, scheduler, external dependency, under-parallelism
High CPU and no progressserialization, compression, bad algorithm, GC
Pods pendingEKS quota, affinity, taint, IPs, node capacity, admission
Serverless job queuedapplication/concurrency capacity or account quota

Adding memory can hide poor partitioning. Adding executors can increase shuffle and source pressure. Tune only after the evidence identifies the bottleneck.

11. Diagnose platform-specific failures

FailureEvidenceSafe diagnosis order
EC2 cluster provisioning failsEMR events, EC2 quota/capacity, subnet IPs, rolesrequest validity, IAM, subnet, capacity, bootstrap
Bootstrap action failsbootstrap stderr in S3, network/DNS, package pathimmutable artifact, checksum, endpoint, command exit
Step remains pendingcluster state, YARN resources, prior stepsconcurrency, capacity, dependency, queue
Serverless job fails before driverjob state details, runtime role, script URIrelease, role trust, S3/KMS, configuration
Serverless workers never reach demandmaximum capacity, quotas, worker shapearithmetic, quota, concurrency, availability
EKS driver pod pendingpod event, quota, node selector, taint, CNIadmission, scheduling, nodes, network
Shuffle fetch failuresexecutor/node history, disk/network, scale eventslost blocks, decommissioning, disk, retry storm
Output is duplicatedrun IDs, retries, sink commitsidempotency and commit protocol
Job succeeds with missing recordscontrol totals, rejected data, partition listschema, filters, late data, silent bad records
Cluster will not terminateprotection, active steps, API eventowner approval, disable exact guard, terminate, verify

Preserve logs before terminating ephemeral compute. Never “fix” a failure by adding broad IAM, opening all egress, disabling encryption, and increasing every capacity limit at once.

12. Select a platform for three workloads

Produce an ADR for each:

Workload A: nightly finance reconciliation

Two-hour deadline, 4 TB Parquet input, strict control totals, JDBC reference lookup, month-end spike, no interactive users, replay-safe output. Compare transient EMR on EC2 fleets with EMR Serverless. Address database connection limits and late capacity.

Workload B: daytime exploratory genomics

Long interactive sessions, specialized native packages, large temporary storage, unpredictable researchers, sensitive data, reproducible environments. Compare long-lived EMR on EC2, EMR Studio, and Serverless interactive choices. Address idle policy and dependency provenance.

Workload C: shared EKS streaming platform

Existing mature EKS team, Kafka input, continuous Spark Structured Streaming, strict namespace quotas, custom image, 30-minute recovery objective. Compare EMR on EKS against Serverless and a dedicated EMR cluster. Address checkpoint compatibility, platform-team on-call, and node upgrades.

Score each option from 1 to 5 for framework fit, startup, elasticity, customization, isolation, networking, reliability, operations, skill, release lifecycle, observability, and total cost. Weight requirements before scoring. A high unweighted total is not an architecture decision.

13. Cost and quota model

EMR on EC2 cost includes EMR charge, EC2 On-Demand/Spot or commitments, EBS, snapshots, S3, public IPv4, load balancers if any, NAT, transfer, logs, metrics, KMS, and idle time. Serverless charges for worker CPU, memory, and storage dimensions plus connected services; pre-initialized capacity adds idle cost. EMR on EKS adds EMR job cost to EKS control plane, nodes/Fargate, autoscaling headroom, storage, observability, image registry, and platform labor.

Model:

monthly cost =
  successful runs
  + failed and retried runs
  + idle or warm capacity
  + storage and requests
  + network and NAT
  + logs, metrics, keys, catalog
  + support and operator time

Include quotas for EC2 instances and vCPUs, Spot, EBS, subnet IPs, EMR clusters, Serverless applications/concurrency/capacity, EKS pods/nodes, ENIs, API rates, Glue partitions, and downstream connections. Capacity unavailable at the deadline is a reliability failure even if average cost is low.

14. Read-only inventory and cleanup proof

aws emr list-clusters \
  --cluster-states STARTING BOOTSTRAPPING RUNNING WAITING \
  --output json | tee active-emr-clusters.json

aws emr-serverless list-applications \
  --output json | tee emr-serverless-applications.json

aws emr-containers list-virtual-clusters \
  --states RUNNING ARRESTED \
  --output json | tee emr-eks-virtual-clusters.json

Inventory alone does not prove zero spend. For every owned experiment verify:

  • EMR on EC2 cluster is TERMINATED or TERMINATED_WITH_ERRORS;
  • related EC2 instances and EBS volumes are gone unless independently owned;
  • Serverless jobs are terminal, applications stopped, and unwanted applications deleted;
  • EMR on EKS jobs are terminal and virtual clusters deleted if lab-owned;
  • lab namespaces, nodes, load balancers, ENIs, logs, S3 outputs, and custom images follow retention policy;
  • Glue, Lake Formation, KMS grants, secrets, and IAM roles are removed only if lab-owned; and
  • delayed billing evidence is reviewed.

Never delete a shared EKS cluster, shared logs bucket, Catalog database, or service-linked role as lesson cleanup.

15. Practical submission

Submit a p18-emr-selection/ package containing:

  1. Spark execution-model explanation and failure boundaries;
  2. pinned release/application dependency matrix and support lifecycle;
  3. three complete workload briefs;
  4. weighted model comparing all three EMR deployments;
  5. EMR on EC2 node, fleet, market, scaling, storage, and termination design;
  6. Serverless application, worker, maximum, concurrency, auto-stop, and runtime-role design;
  7. EKS namespace, RBAC, execution-role, scheduler, image, and owner design;
  8. S3, Catalog, Lake Formation, KMS, and table-format map;
  9. dependency-by-dependency network path;
  10. retry, Spot, idempotency, checkpoint, and output-commit strategy;
  11. baseline and optimized Spark evidence from supplied metrics or an approved run;
  12. diagnosis of one failure from each deployment model;
  13. monthly cost model with failure and idle sensitivity;
  14. quota and capacity-risk register;
  15. operational runbook and escalation map; and
  16. termination, retained-state, and delayed billing evidence.

Knowledge check

  1. Why is an EMR release label more than a Spark version? It pins a tested software and patch combination.
  2. Which EC2 nodes are usually safest for Spot? Replayable task capacity, not critical primary state.
  3. Does termination protection preserve HDFS? No; it is an accidental-termination guard, not backup.
  4. What can Serverless pre-initialized capacity cost? It remains billed while ready in a started application.
  5. What bounds a Serverless application's scale? Maximum CPU, memory, and disk plus quotas and worker shapes.
  6. Does a virtual cluster create EKS? No; it registers an EKS namespace with EMR.
  7. Why can an EKS job stay pending after EMR accepts it? Kubernetes scheduling, admission, nodes, quota, or CNI capacity can block pods.
  8. Why can more executors make a job worse? They can increase shuffle, overhead, and downstream pressure.
  9. What makes retry safe? Idempotent or transactional output, durable checkpoints, and reconciliation.
  10. What proves cleanup? Terminal compute plus absence or governed retention of every connected billable resource.

Lesson acceptance

Pass only when all are true:

  • Spark driver/executor/task/shuffle behavior is explained before platform selection;
  • exact release labels, dependencies, Region availability, and support lifecycle are recorded;
  • all three EMR deployment models are compared against the same requirements;
  • EC2 primary/core/task, group/fleet, On-Demand/Spot, scaling, storage, and termination decisions are justified;
  • Serverless role, worker shape, maximum, concurrency, warm capacity, network, and auto-stop are bounded;
  • EKS virtual cluster, namespace, RBAC, IAM, scheduling, image, network, and ownership are complete;
  • durable state is outside ephemeral compute and recovery semantics are testable;
  • identity, encryption, KMS, Lake Formation, and every private network dependency are mapped;
  • failure diagnosis uses logs, metrics, events, and output behavior before resizing;
  • retries, interruption, partial output, and streaming checkpoints have explicit safety;
  • three workload ADRs use weighted evidence and reject unsuitable options;
  • cost includes idle, failure, storage, network, logs, keys, platform, and labor;
  • quotas and capacity-at-deadline risk are recorded; and
  • cleanup proves compute termination and handles all connected resources.

Official sources

Advertisement