All infographics

DATA PROCESSING / 09

Amazon EMR

Process large datasets with managed analytics engines

Download sheet SVG

THE BIG PICTURE

Turn datasets into results

  1. Source dataS3 / external systems
  2. EMR jobSpark / supported frameworks
  3. ResultsS3 / downstream analytics

Example data pipeline: an EMR job reads source data and writes results. Choose EC2, Serverless, or EKS to match your operational needs.

01Three deployment choices

  • Amazon EMR runs open-source analytics engines such as Spark and Hive for data processing
  • EMR on EC2 provisions clusters in your account, with control over instances and installed applications
  • EMR Serverless runs Spark or Hive jobs without requiring you to provision a cluster
  • EMR on EKS runs analytics jobs as containers on an EKS cluster that your team manages
  • Engine versions and features vary by release and deployment option; check compatibility first

02EC2 primary, core and task nodes

  • The primary node coordinates cluster services and monitors workload and node health
  • Core nodes run computation and hold HDFS data; removing them carelessly can lose local data
  • Task nodes add compute capacity and do not hold HDFS data, making them useful for elastic workers
  • Instance groups use one instance type per group; fleets can mix types and On-Demand with Spot
  • Managed scaling adjusts supported clusters within limits you set; protect essential capacity from churn

03EMR Serverless applications

  • An application selects an EMR release and engine, such as Spark or Hive, before jobs are submitted
  • Each job run uses a runtime IAM role for access to its input data, output locations and other resources
  • Workers scale with job demand; maximum application capacity bounds how many resources may be used
  • Pre-initialized workers reduce startup delay but incur charges while that capacity remains running
  • Use separate applications for teams or environments that need independent versions and cost controls

04EMR on EKS virtual clusters

  • A virtual cluster registers one Kubernetes namespace on an existing Amazon EKS cluster
  • Submit a Spark JAR, PySpark script or Spark SQL job with an execution role and EMR release
  • EMR supplies analytics containers while Kubernetes schedules their driver and executor pods
  • Your team still manages EKS capacity, networking and namespace access for shared workloads
  • Different virtual clusters can share an EKS cluster; registering one does not provision new compute

05Releases and dependencies

  • An EMR release label selects a tested bundle of engine and dependency versions
  • Select the applications you need on EC2; Serverless applications choose a single engine type
  • Pin Python packages, JARs and connectors to compatible Spark, Scala and runtime versions
  • Use bootstrap actions for EC2 node setup; supported Serverless and EKS custom images package dependencies
  • Test jobs and table formats against a new release before replacing production environments

06S3, HDFS and local storage

  • Keep durable inputs and outputs in S3 so they survive an EC2 cluster being terminated
  • HDFS uses cluster storage and provides data locality; its data does not survive cluster termination
  • S3A is the default S3 connector in newer releases; older releases may use EMRFS
  • Local shuffle and spill files can fill worker disks; size or select storage for the workload
  • Use columnar formats and useful partitioning to reduce data scanned by analytics jobs

07Tables and catalogs

  • A catalog stores table metadata; the underlying rows may live in S3 independently of that catalog
  • Configure AWS Glue Data Catalog as a persistent Hive metastore or a supported Iceberg catalog
  • Iceberg, Hudi and Delta Lake capabilities depend on the EMR release, engine and deployment option
  • Keep table data in an explicit S3 location when multiple clusters or query engines need it
  • Lake Formation integrations can govern lake data; verify the supported access mode and release

08Spark sizing and performance

  • The driver coordinates a Spark application; executors run its distributed tasks
  • Size driver and executor memory separately, including runtime overhead in worker memory requirements
  • Joins and aggregations can shuffle data across workers, consuming network, memory and disk
  • Skewed keys create slow tasks; inspect partitions and use suitable join and adaptive-query settings
  • Avoid collecting a large distributed dataset into the driver; write results or sample them instead

09Jobs and pipeline orchestration

  • On EC2, submit processing steps to a cluster; configure concurrency and failure actions deliberately
  • Serverless and EMR on EKS expose job-run APIs with entry points, arguments and job configuration
  • Step Functions or Airflow can coordinate dependent jobs, retries and cleanup around your pipeline
  • Store versioned job code and explicit input and output paths so a run can be reproduced
  • Retries should use safe output and commit logic; replaying a job must not silently duplicate data

10Security and networking

  • EC2 service roles manage infrastructure; instance profiles and supported runtime roles grant data access
  • Serverless and EKS job execution roles need scoped permissions for S3, catalogs and any KMS keys
  • On EC2, a security configuration can enable encryption for local data and traffic between services
  • Private deployments need routes or endpoints to their data sources and required AWS APIs
  • Restrict cluster and notebook access; keep credentials out of bootstrap scripts and job arguments

11Logs and troubleshooting

  • Use Spark UI or history views to inspect stages, task timings, executor failures and shuffle spill
  • Driver and executor logs answer different questions; retain both when investigating job failures
  • Serverless supports managed log storage, S3 and CloudWatch; choose logging options before a run
  • CloudWatch metrics expose workload and resource usage; alert on failed jobs and prolonged queues
  • Preserve useful logs and event history before terminating an EC2 cluster or deleting test resources

12Costs and operational pitfalls

  • EMR on EC2 adds an EMR charge to EC2 and any EBS costs; idle clusters still consume resources
  • Serverless charges for worker vCPU, memory and applicable storage while workers are running
  • EMR on EKS adds an EMR charge to EKS and its underlying compute and storage costs
  • S3 requests, logs, catalogs, networking and other connected services can add separate charges
  • On EC2, use suitable Spot task nodes and auto-termination; plan for interruptions and startup delay

Go to the source

Use AWS documentation for current limits, availability, and pricing.

EMR clusters, nodes and deployment concepts EMR Serverless concepts EMR on EKS concepts EMR release bundles and application versions Security in Amazon EMR EMR pricing by deployment option