DATA PROCESSING / 09
Amazon EMR
Process large datasets with managed analytics engines
THE BIG PICTURE
Turn datasets into results
- Source dataS3 / external systems
- EMR jobSpark / supported frameworks
- 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.