{"id":4635,"date":"2026-07-23T14:15:48","date_gmt":"2026-07-23T14:15:48","guid":{"rendered":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/2026\/07\/23\/building-reliable-emr-pipelines-with-custom-amis-and-step-functions\/"},"modified":"2026-07-23T14:15:48","modified_gmt":"2026-07-23T14:15:48","slug":"building-reliable-emr-pipelines-with-custom-amis-and-step-functions","status":"publish","type":"post","link":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/2026\/07\/23\/building-reliable-emr-pipelines-with-custom-amis-and-step-functions\/","title":{"rendered":"Building Reliable EMR Pipelines With Custom AMIs and Step Functions"},"content":{"rendered":"<div><img data-opt-id=523763747  fetchpriority=\"high\" decoding=\"async\" width=\"770\" height=\"330\" src=\"https:\/\/devops.com\/wp-content\/uploads\/2021\/05\/securecoding.jpg\" class=\"attachment-large size-large wp-post-image\" alt=\"OWASP DevSecOps\" \/><\/div>\n<p><img data-opt-id=504193142  fetchpriority=\"high\" decoding=\"async\" width=\"150\" height=\"150\" src=\"https:\/\/devops.com\/wp-content\/uploads\/2021\/05\/securecoding-150x150.jpg\" class=\"attachment-thumbnail size-thumbnail wp-post-image\" alt=\"OWASP DevSecOps\" \/><\/p>\n<p>A practical pattern for replacing runtime bootstrap fragility with image-based dependencies, orchestration and retry-aware recovery.<\/p>\n<h3>The Problem: Dependency Drift in EMR Pipelines<\/h3>\n<p>Amazon EMR remains a strong option for running Apache Spark workloads, but many production-style pipelines still rely on bootstrap actions to install dependencies when a cluster starts. That works for early experiments. It becomes fragile when pipelines need predictable startup behavior, repeatable dependency versions and cleaner recovery from transient infrastructure failures.<\/p>\n<p>The common failure mode is simple: Every cluster launch repeats package installation, network calls and setup scripts. If a package repository is slow, a dependency resolver changes behavior or a security patch needs to be rolled across many jobs, the pipeline team is debugging environment creation instead of data logic.<\/p>\n<p>One way to reduce that operational noise is to treat the EMR runtime as immutable infrastructure: Custom AMIs with pre-baked dependencies, automated image builds, versioned metadata, Step Functions orchestration and failure handling that understands retryable conditions.<\/p>\n<h3>Why Custom AMIs Beat Bootstrap Actions<\/h3>\n<h3>Architecture Overview<\/h3>\n<p>\u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510 \u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510 \u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510<br \/>\n\u2502 CodeCommit \u2502\u2500\u2500\u2500\u2500<img data-opt-id=215683849  data-opt-src=\"https:\/\/s.w.org\/images\/core\/emoji\/16.0.1\/72x72\/25b6.png\"  decoding=\"async\" src=\"data:image/svg+xml,%3Csvg%20viewBox%3D%220%200%20100%%20100%%22%20width%3D%22100%%22%20height%3D%22100%%22%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%3E%3Crect%20width%3D%22100%%22%20height%3D%22100%%22%20fill%3D%22transparent%22%2F%3E%3C%2Fsvg%3E\" alt=\"\u25b6\" class=\"wp-smiley\" \/>\u2502 Lambda Builder \u2502\u2500\u2500\u2500\u2500<img data-opt-id=215683849  data-opt-src=\"https:\/\/s.w.org\/images\/core\/emoji\/16.0.1\/72x72\/25b6.png\"  decoding=\"async\" src=\"data:image/svg+xml,%3Csvg%20viewBox%3D%220%200%20100%%20100%%22%20width%3D%22100%%22%20height%3D%22100%%22%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%3E%3Crect%20width%3D%22100%%22%20height%3D%22100%%22%20fill%3D%22transparent%22%2F%3E%3C%2Fsvg%3E\" alt=\"\u25b6\" class=\"wp-smiley\" \/>\u2502 EC2 (Packer) \u2502<br \/>\n\u2502 (requirements) \u2502 \u2502 (trigger AMI \u2502 \u2502 (Custom AMI) \u2502<br \/>\n\u2502 \u2502 \u2502 build) \u2502 \u2502 \u2502<br \/>\n\u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518 \u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518 \u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u252c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<br \/>\n\u2502<br \/>\n\u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<br \/>\n\u25bc<br \/>\n\u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510<br \/>\n\u2502 Golden AMI \u2502<br \/>\n\u2502 (EMR + Python) \u2502<br \/>\n\u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u252c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<br \/>\n\u2502<br \/>\n\u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510 \u2502 \u250c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2510<br \/>\n\u2502 Step Functions \u2502<img data-opt-id=1738152803  data-opt-src=\"https:\/\/s.w.org\/images\/core\/emoji\/16.0.1\/72x72\/25c0.png\"  decoding=\"async\" src=\"data:image/svg+xml,%3Csvg%20viewBox%3D%220%200%20100%%20100%%22%20width%3D%22100%%22%20height%3D%22100%%22%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%3E%3Crect%20width%3D%22100%%22%20height%3D%22100%%22%20fill%3D%22transparent%22%2F%3E%3C%2Fsvg%3E\" alt=\"\u25c0\" class=\"wp-smiley\" \/>\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2534\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500<img data-opt-id=215683849  data-opt-src=\"https:\/\/s.w.org\/images\/core\/emoji\/16.0.1\/72x72\/25b6.png\"  decoding=\"async\" src=\"data:image/svg+xml,%3Csvg%20viewBox%3D%220%200%20100%%20100%%22%20width%3D%22100%%22%20height%3D%22100%%22%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%3E%3Crect%20width%3D%22100%%22%20height%3D%22100%%22%20fill%3D%22transparent%22%2F%3E%3C%2Fsvg%3E\" alt=\"\u25b6\" class=\"wp-smiley\" \/>\u2502 EMR Cluster \u2502<br \/>\n\u2502 Orchestrator \u2502 \u2502 (Custom AMI) \u2502<br \/>\n\u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u252c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518 \u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u252c\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<br \/>\n\u2502 \u2502<br \/>\n\u2502 Error? \u2502<br \/>\n\u2502<img data-opt-id=1738152803  data-opt-src=\"https:\/\/s.w.org\/images\/core\/emoji\/16.0.1\/72x72\/25c0.png\"  decoding=\"async\" src=\"data:image/svg+xml,%3Csvg%20viewBox%3D%220%200%20100%%20100%%22%20width%3D%22100%%22%20height%3D%22100%%22%20xmlns%3D%22http%3A%2F%2Fwww.w3.org%2F2000%2Fsvg%22%3E%3Crect%20width%3D%22100%%22%20height%3D%22100%%22%20fill%3D%22transparent%22%2F%3E%3C%2Fsvg%3E\" alt=\"\u25c0\" class=\"wp-smiley\" \/>\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500 Retry logic \u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<br \/>\n\u2502<br \/>\n\u250c\u2500\u2500\u2500\u2500\u25bc\u2500\u2500\u2500\u2500\u2510<br \/>\n\u2502 DLQ \u2502 (CloudWatch analysis)<br \/>\n\u2514\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2500\u2518<\/p>\n<h3>Part 1: Building the Custom AMI<\/h3>\n<h3>Step 1: Create the Packer Configuration<\/h3>\n<p>Use HashiCorp Packer to automate AMI builds. This creates a reproducible, version-controlled image pipeline.<\/p>\n<p># emr-custom-ami.pkr.hcl<br \/>\npacker {<br \/>\nrequired_plugins {<br \/>\namazon = {<br \/>\nsource = \u201cgithub.com\/hashicorp\/amazon\u201d<br \/>\nversion = \u201c~&gt; 1\u201d<br \/>\n}<br \/>\n}<br \/>\n}<\/p>\n<p>variable \u201cemr_release\u201d {<br \/>\ndefault = \u201cemr-6.15.0\u201d<br \/>\n}<\/p>\n<p>variable \u201cpython_packages\u201d {<br \/>\ndefault = [<br \/>\n\u201cpandas==2.0.3\u201d,<br \/>\n\u201cnumpy==1.24.3\u201d,<br \/>\n\u201cboto3==1.28.25\u201d,<br \/>\n\u201cpyspark==3.4.1\u201d,<br \/>\n\u201cgreat-expectations==0.17.0\u201d<br \/>\n]<br \/>\n}<\/p>\n<p>source \u201camazon-ebs\u201d \u201cemr-custom\u201d {<br \/>\nami_name = \u201cemr-${var.emr_release}-analytics-${formatdate(\u201dYYYYMMDD-hhmm\u201d, timestamp())}\u201d<br \/>\ninstance_type = \u201cm5.xlarge\u201d<br \/>\nregion = \u201cus-east-1\u201d<\/p>\n<p># Start from official EMR AMI<br \/>\nsource_ami_filter {<br \/>\nfilters = {<br \/>\nname = \u201camzn2-ami-emr-${var.emr_release}*\u201d<br \/>\nroot-device-type = \u201cebs\u201d<br \/>\nvirtualization-type = \u201chvm\u201d<br \/>\n}<br \/>\nowners = [\u201camazon\u201d]<br \/>\nmost_recent = true<br \/>\n}<\/p>\n<p>ssh_username = \u201cec2-user\u201d<\/p>\n<p># EBS optimization for Spark<br \/>\nlaunch_block_device_mappings {<br \/>\ndevice_name = \u201c\/dev\/xvda\u201d<br \/>\nvolume_size = 100<br \/>\nvolume_type = \u201cgp3\u201d<br \/>\ndelete_on_termination = true<br \/>\n}<\/p>\n<p>tags = {<br \/>\nName = \u201cEMR Analytics AMI\u201d<br \/>\nEMRRelease = var.emr_release<br \/>\nBuiltBy = \u201cPacker\u201d<br \/>\nRepository = \u201cgithub.com\/example-org\/emr-ami-pipeline\u201d<br \/>\n}<br \/>\n}<\/p>\n<p>build {<br \/>\nsources = [\u201csource.amazon-ebs.emr-custom\u201d]\n<\/p>\n<p># Install Python 3.9 (EMR default is often 3.7)<br \/>\nprovisioner \u201cshell\u201d {<br \/>\ninline = [<br \/>\n\u201csudo amazon-linux-extras install python3.9 -y\u201d,<br \/>\n\u201csudo alternatives \u2013install \/usr\/bin\/python3 python3 \/usr\/bin\/python3.9 1\u201d,<br \/>\n\u201csudo python3.9 -m pip install \u2013upgrade pip setuptools wheel\u201d<br \/>\n]<br \/>\n}<\/p>\n<p># Install Python dependencies<br \/>\nprovisioner \u201cshell\u201d {<br \/>\ninline = concat(<br \/>\n[\u201csudo python3.9 -m pip install \u201c],<br \/>\n[join(\u201d \u201c, var.python_packages)]<br \/>\n)<br \/>\n}<\/p>\n<p># Verify installations<br \/>\nprovisioner \u201cshell\u201d {<br \/>\ninline = [<br \/>\n\u201cpython3 \u2013version\u201d,<br \/>\n\u201cpip3 list | grep pandas\u201d,<br \/>\n\u201cpip3 list | grep numpy\u201d<br \/>\n]<br \/>\n}<\/p>\n<p># Clean up for AMI optimization<br \/>\nprovisioner \u201cshell\u201d {<br \/>\ninline = [<br \/>\n\u201csudo yum clean all\u201d,<br \/>\n\u201csudo rm -rf \/var\/cache\/yum\u201d,<br \/>\n\u201csudo rm -f \/home\/ec2-user\/.ssh\/authorized_keys\u201d,<br \/>\n\u201csudo rm -f \/root\/.ssh\/authorized_keys\u201d<br \/>\n]<br \/>\n}<br \/>\n}<\/p>\n<h3>Step 2: Lambda-Powered AMI Builder<\/h3>\n<p>Trigger AMI builds automatically when requirements.txt changes:<\/p>\n<p># lambda\/ami_builder.py<br \/>\nimport json<br \/>\nimport boto3<br \/>\nimport os<\/p>\n<p>ec2 = boto3.client(\u2018ec2\u2019)<br \/>\nssm = boto3.client(\u2018ssm\u2019)<\/p>\n<p>def lambda_handler(event, context):<br \/>\n\u201c\u201d\u201d<br \/>\nTriggered by CodeCommit push to requirements.txt<br \/>\nBuilds new EMR custom AMI using Packer<br \/>\n\u201c\u201d\u201d<\/p>\n<p># Get latest EMR release from SSM Parameter Store<br \/>\nemr_release = ssm.get_parameter(<br \/>\nName=\u2019\/emr-pipeline\/emr-release\u2019<br \/>\n)[\u2018Parameter\u2019][\u2018Value\u2019]\n<\/p>\n<p># Start Packer build on EC2 instance (builder pattern)<br \/>\nuser_data = f\u201d\u2019#!\/bin\/bash<br \/>\nyum update -y<br \/>\nyum install -y packer git<\/p>\n<p># Clone repo with Packer templates<br \/>\ncd \/tmp<br \/>\ngit clone https:\/\/github.com\/example-org\/emr-ami-pipeline.git<br \/>\ncd emr-ami-pipeline<\/p>\n<p># Run Packer build<br \/>\npacker init emr-custom-ami.pkr.hcl<br \/>\npacker build <br \/>\n-var \u2019emr_release={emr_release}\u2019 <br \/>\n-var \u2018python_packages_file=requirements.txt\u2019 <br \/>\nemr-custom-ami.pkr.hcl<\/p>\n<p># Notify Step Functions\/SNS on completion<br \/>\naws sns publish <br \/>\n\u2013topic-arn {os.environ[\u2018SNS_TOPIC_ARN\u2019]} <br \/>\n\u2013message \u201cAMI build completed\u201d<br \/>\n\u201d\u2019<\/p>\n<p># Launch builder instance<br \/>\nresponse = ec2.run_instances(<br \/>\nImageId=\u2019ami-0c55b159cbfafe1f0\u2032, # Amazon Linux 2<br \/>\nInstanceType=\u2019m5.2xlarge\u2019,<br \/>\nMinCount=1,<br \/>\nMaxCount=1,<br \/>\nIamInstanceProfile={\u2018Name\u2019: \u2018EMR-AMI-Builder\u2019},<br \/>\nUserData=user_data,<br \/>\nTagSpecifications=[{<br \/>\n\u2018ResourceType\u2019: \u2018instance\u2019,<br \/>\n\u2018Tags\u2019: [<br \/>\n{\u2018Key\u2019: \u2018Name\u2019, \u2018Value\u2019: \u2019emr-ami-builder\u2019},<br \/>\n{\u2018Key\u2019: \u2018AutoTerminate\u2019, \u2018Value\u2019: \u2018true\u2019}<br \/>\n]<br \/>\n}]<br \/>\n)<\/p>\n<p>return {<br \/>\n\u2018statusCode\u2019: 200,<br \/>\n\u2018body\u2019: json.dumps({<br \/>\n\u2018message\u2019: \u2018AMI build initiated\u2019,<br \/>\n\u2018instanceId\u2019: response[\u2018Instances\u2019][0][\u2018InstanceId\u2019]<br \/>\n})<br \/>\n}<\/p>\n<h3>Step 3: Store Golden AMI Metadata<\/h3>\n<p>Track AMI versions in DynamoDB for pipeline lookups:<\/p>\n<p># lambda\/ami_registry.py<br \/>\nimport boto3<br \/>\nfrom datetime import datetime<\/p>\n<p>dynamodb = boto3.resource(\u2018dynamodb\u2019)<br \/>\ntable = dynamodb.Table(\u2019emr-golden-amis\u2019)<\/p>\n<p>def register_ami(ami_id, emr_release, packages, commit_hash):<br \/>\n\u201c\u201d\u201dRegister newly built AMI as \u2018active\u2019 for Step Functions.\u201d\u201d\u201d<\/p>\n<p># Mark previous AMI as deprecated<br \/>\ntable.update_item(<br \/>\nKey={\u2019emr_release\u2019: emr_release, \u2018status\u2019: \u2018active\u2019},<br \/>\nUpdateExpression=\u2019SET #status = :deprecated\u2019,<br \/>\nExpressionAttributeNames={\u2018#status\u2019: \u2018status\u2019},<br \/>\nExpressionAttributeValues={\u2018:deprecated\u2019: \u2018deprecated\u2019}<br \/>\n)<\/p>\n<p># Register new AMI<br \/>\ntable.put_item(Item={<br \/>\n\u2018ami_id\u2019: ami_id,<br \/>\n\u2019emr_release\u2019: emr_release,<br \/>\n\u2018status\u2019: \u2018active\u2019,<br \/>\n\u2018python_packages\u2019: packages,<br \/>\n\u2018commit_hash\u2019: commit_hash,<br \/>\n\u2018created_at\u2019: datetime.utcnow().isoformat(),<br \/>\n\u2018ttl\u2019: int((datetime.utcnow().timestamp()) + (90 * 86400)) # 90-day retention<br \/>\n})<\/p>\n<h3>Part 2: Step Functions Orchestration<\/h3>\n<h3>State Machine Definition<\/h3>\n<p>{<br \/>\n\u201cComment\u201d: \u201cEMR Data Pipeline with Custom AMI and Error Handling\u201d,<br \/>\n\u201cStartAt\u201d: \u201cGetGoldenAMI\u201d,<br \/>\n\u201cStates\u201d: {<br \/>\n\u201cGetGoldenAMI\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::lambda:invoke\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cFunctionName\u201d: \u201cget-golden-ami\u201d,<br \/>\n\u201cPayload\u201d: {<br \/>\n\u201cemr_release\u201d: \u201cemr-6.15.0\u201d<br \/>\n}<br \/>\n},<br \/>\n\u201cResultPath\u201d: \u201c$.ami\u201d,<br \/>\n\u201cNext\u201d: \u201cCreateEMRCluster\u201d<br \/>\n},<\/p>\n<p>\u201cCreateEMRCluster\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::elasticmapreduce:createCluster.sync\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cName\u201d: \u201cpipeline-${$.executionId}\u201d,<br \/>\n\u201cReleaseLabel\u201d: \u201cemr-6.15.0\u201d,<br \/>\n\u201cCustomAmiId.$\u201d: \u201c$.ami.Payload.ami_id\u201d,<br \/>\n\u201cLogUri\u201d: \u201cs3:\/\/emr-logs-bucket\/pipelines\/\u201d,<br \/>\n\u201cInstances\u201d: {<br \/>\n\u201cKeepJobFlowAliveWhenNoSteps\u201d: false,<br \/>\n\u201cInstanceFleets\u201d: [<br \/>\n{<br \/>\n\u201cName\u201d: \u201cMaster\u201d,<br \/>\n\u201cInstanceFleetType\u201d: \u201cMASTER\u201d,<br \/>\n\u201cTargetOnDemandCapacity\u201d: 1,<br \/>\n\u201cInstanceTypeConfigs\u201d: [<br \/>\n{\u201cInstanceType\u201d: \u201cm5.xlarge\u201d}<br \/>\n]<br \/>\n},<br \/>\n{<br \/>\n\u201cName\u201d: \u201cCore\u201d,<br \/>\n\u201cInstanceFleetType\u201d: \u201cCORE\u201d,<br \/>\n\u201cTargetSpotCapacity\u201d: 4,<br \/>\n\u201cInstanceTypeConfigs\u201d: [<br \/>\n{\u201cInstanceType\u201d: \u201cm5.2xlarge\u201d, \u201cWeightedCapacity\u201d: 2},<br \/>\n{\u201cInstanceType\u201d: \u201cm5.4xlarge\u201d, \u201cWeightedCapacity\u201d: 4}<br \/>\n]<br \/>\n}<br \/>\n],<br \/>\n\u201cEc2SubnetIds\u201d: [\u201csubnet-12345abcde\u201d],<br \/>\n\u201cEmrManagedMasterSecurityGroup\u201d: \u201csg-master\u201d,<br \/>\n\u201cEmrManagedSlaveSecurityGroup\u201d: \u201csg-core\u201d<br \/>\n},<br \/>\n\u201cApplications\u201d: [<br \/>\n{\u201cName\u201d: \u201cSpark\u201d},<br \/>\n{\u201cName\u201d: \u201cHadoop\u201d}<br \/>\n],<br \/>\n\u201cServiceRole\u201d: \u201cEMR_DefaultRole\u201d,<br \/>\n\u201cJobFlowRole\u201d: \u201cEMR_EC2_DefaultRole\u201d<br \/>\n},<br \/>\n\u201cResultPath\u201d: \u201c$.cluster\u201d,<br \/>\n\u201cNext\u201d: \u201cSubmitSparkJob\u201d<br \/>\n},<\/p>\n<p>\u201cSubmitSparkJob\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::elasticmapreduce:addStep.sync\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cClusterId.$\u201d: \u201c$.cluster.ClusterId\u201d,<br \/>\n\u201cStep\u201d: {<br \/>\n\u201cName\u201d: \u201cDataProcessing\u201d,<br \/>\n\u201cActionOnFailure\u201d: \u201cCONTINUE\u201d,<br \/>\n\u201cHadoopJarStep\u201d: {<br \/>\n\u201cJar\u201d: \u201ccommand-runner.jar\u201d,<br \/>\n\u201cArgs\u201d: [<br \/>\n\u201cspark-submit\u201d,<br \/>\n\u201c\u2013deploy-mode\u201d, \u201ccluster\u201d,<br \/>\n\u201c\u2013conf\u201d, \u201cspark.dynamicAllocation.enabled=true\u201d,<br \/>\n\u201cs3:\/\/emr-jobs\/process_data.py\u201d,<br \/>\n\u201c\u2013input\u201d, \u201cs3:\/\/data-bucket\/raw\/\u201d,<br \/>\n\u201c\u2013output\u201d, \u201cs3:\/\/data-bucket\/processed\/\u201d<br \/>\n]<br \/>\n}<br \/>\n}<br \/>\n},<br \/>\n\u201cResultPath\u201d: \u201c$.step\u201d,<br \/>\n\u201cNext\u201d: \u201cCheckStepStatus\u201d<br \/>\n},<\/p>\n<p>\u201cCheckStepStatus\u201d: {<br \/>\n\u201cType\u201d: \u201cChoice\u201d,<br \/>\n\u201cChoices\u201d: [<br \/>\n{<br \/>\n\u201cVariable\u201d: \u201c$.step.Step.State\u201d,<br \/>\n\u201cStringEquals\u201d: \u201cCOMPLETED\u201d,<br \/>\n\u201cNext\u201d: \u201cSuccess\u201d<br \/>\n},<br \/>\n{<br \/>\n\u201cVariable\u201d: \u201c$.step.Step.State\u201d,<br \/>\n\u201cStringEquals\u201d: \u201cFAILED\u201d,<br \/>\n\u201cNext\u201d: \u201cAnalyzeFailure\u201d<br \/>\n}<br \/>\n],<br \/>\n\u201cDefault\u201d: \u201cWaitForCompletion\u201d<br \/>\n},<\/p>\n<p>\u201cWaitForCompletion\u201d: {<br \/>\n\u201cType\u201d: \u201cWait\u201d,<br \/>\n\u201cSeconds\u201d: 60,<br \/>\n\u201cNext\u201d: \u201cCheckStepStatus\u201d<br \/>\n},<\/p>\n<p>\u201cAnalyzeFailure\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::lambda:invoke\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cFunctionName\u201d: \u201canalyze-emr-failure\u201d,<br \/>\n\u201cPayload\u201d: {<br \/>\n\u201ccluster_id.$\u201d: \u201c$.cluster.ClusterId\u201d,<br \/>\n\u201cstep_id.$\u201d: \u201c$.step.Step.Id\u201d,<br \/>\n\u201cattempt_count.$\u201d: \u201c$.attempt_count\u201d<br \/>\n}<br \/>\n},<br \/>\n\u201cResultPath\u201d: \u201c$.failure_analysis\u201d,<br \/>\n\u201cNext\u201d: \u201cShouldRetry\u201d<br \/>\n},<\/p>\n<p>\u201cShouldRetry\u201d: {<br \/>\n\u201cType\u201d: \u201cChoice\u201d,<br \/>\n\u201cChoices\u201d: [<br \/>\n{<br \/>\n\u201cAnd\u201d: [<br \/>\n{\u201cVariable\u201d: \u201c$.failure_analysis.Payload.retryable\u201d, \u201cBooleanEquals\u201d: true},<br \/>\n{\u201cVariable\u201d: \u201c$.attempt_count\u201d, \u201cNumericLessThan\u201d: 3}<br \/>\n],<br \/>\n\u201cNext\u201d: \u201cCalculateBackoff\u201d<br \/>\n}<br \/>\n],<br \/>\n\u201cDefault\u201d: \u201cSendToDLQ\u201d<br \/>\n},<\/p>\n<p>\u201cCalculateBackoff\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::lambda:invoke\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cFunctionName\u201d: \u201ccalculate-backoff\u201d,<br \/>\n\u201cPayload\u201d: {<br \/>\n\u201cattempt.$\u201d: \u201c$.attempt_count\u201d<br \/>\n}<br \/>\n},<br \/>\n\u201cResultPath\u201d: \u201c$.backoff\u201d,<br \/>\n\u201cNext\u201d: \u201cWaitBeforeRetry\u201d<br \/>\n},<\/p>\n<p>\u201cWaitBeforeRetry\u201d: {<br \/>\n\u201cType\u201d: \u201cWait\u201d,<br \/>\n\u201cSecondsPath\u201d: \u201c$.backoff.Payload.seconds\u201d,<br \/>\n\u201cNext\u201d: \u201cCreateEMRCluster\u201d<br \/>\n},<\/p>\n<p>\u201cSendToDLQ\u201d: {<br \/>\n\u201cType\u201d: \u201cTask\u201d,<br \/>\n\u201cResource\u201d: \u201carn:aws:states:::sqs:sendMessage\u201d,<br \/>\n\u201cParameters\u201d: {<br \/>\n\u201cQueueUrl\u201d: \u201chttps:\/\/sqs.us-east-1.amazonaws.com\/123456789\/emr-failures-dlq\u201d,<br \/>\n\u201cMessageBody\u201d: {<br \/>\n\u201cexecution_id.$\u201d: \u201c$$.Execution.Id\u201d,<br \/>\n\u201ccluster_id.$\u201d: \u201c$.cluster.ClusterId\u201d,<br \/>\n\u201cerror.$\u201d: \u201c$.failure_analysis.Payload.error_type\u201d,<br \/>\n\u201ctimestamp.$\u201d: \u201c$$.State.EnteredTime\u201d<br \/>\n}<br \/>\n},<br \/>\n\u201cNext\u201d: \u201cPipelineFailed\u201d<br \/>\n},<\/p>\n<p>\u201cSuccess\u201d: {<br \/>\n\u201cType\u201d: \u201cSucceed\u201d<br \/>\n},<\/p>\n<p>\u201cPipelineFailed\u201d: {<br \/>\n\u201cType\u201d: \u201cFail\u201d,<br \/>\n\u201cError\u201d: \u201cPipelineFailed\u201d,<br \/>\n\u201cCause\u201d: \u201cNon-retryable error or max retries exceeded\u201d<br \/>\n}<br \/>\n}<br \/>\n}<\/p>\n<h3>Error Analysis Lambda<\/h3>\n<p># lambda\/analyze_emr_failure.py<br \/>\nimport boto3<br \/>\nimport json<\/p>\n<p>emr = boto3.client(\u2019emr\u2019)<br \/>\nlogs = boto3.client(\u2018logs\u2019)<\/p>\n<p>def lambda_handler(event, context):<br \/>\ncluster_id = event[\u2018cluster_id\u2019]<br \/>\nstep_id = event[\u2018step_id\u2019]<br \/>\nattempt = event.get(\u2018attempt_count\u2019, 0)<\/p>\n<p># Get step details<br \/>\nresponse = emr.describe_step(ClusterId=cluster_id, StepId=step_id)<br \/>\nstep_status = response[\u2018Step\u2019][\u2018Status\u2019]\n<\/p>\n<p>error_info = step_status.get(\u2018FailureDetails\u2019, {})<br \/>\nreason = error_info.get(\u2018Reason\u2019, \u2018Unknown\u2019)<br \/>\nlog_file = error_info.get(\u2018LogFile\u2019, \u201d)<\/p>\n<p># Classify error type<br \/>\nretryable_errors = [<br \/>\n\u2018Spot instance termination\u2019,<br \/>\n\u2018Container killed on request\u2019,<br \/>\n\u2018Connection timeout\u2019,<br \/>\n\u2018ThrottlingException\u2019<br \/>\n]\n<\/p>\n<p>is_retryable = any(err in reason for err in retryable_errors)<\/p>\n<p># If spot termination, suggest on-demand retry<br \/>\nif \u2018Spot\u2019 in reason and attempt == 2:<br \/>\nis_retryable = False # Third failure, use DLQ<\/p>\n<p>return {<br \/>\n\u2018retryable\u2019: is_retryable,<br \/>\n\u2018error_type\u2019: reason,<br \/>\n\u2018log_file\u2019: log_file,<br \/>\n\u2018recommendation\u2019: \u2018retry_on_demand\u2019 if \u2018Spot\u2019 in reason else \u2018retry\u2019<br \/>\n}<\/p>\n<p>def calculate_backoff(attempt):<br \/>\n\u201c\u201d\u201dExponential backoff: 30s, 2min, 5min\u201d\u201d\u201d<br \/>\ndelays = [30, 120, 300]<br \/>\nreturn {\u2018seconds\u2019: delays[min(attempt, 2)]}<\/p>\n<h3>Part 3: Production Hardening<\/h3>\n<h3>Spot Instance Handling With Checkpointing<\/h3>\n<p># spark_job.py (runs on EMR)<br \/>\nfrom pyspark.sql import SparkSession<br \/>\nimport boto3<br \/>\nimport sys<\/p>\n<p>spark = SparkSession.builder <br \/>\n.appName(\u201cDataPipeline\u201d) <br \/>\n.config(\u201cspark.sql.adaptive.enabled\u201d, \u201ctrue\u201d) <br \/>\n.config(\u201cspark.checkpoint.dir\u201d, \u201cs3:\/\/emr-checkpoints\/job\/\u201d) <br \/>\n.getOrCreate()<\/p>\n<p>def main():<br \/>\n# Check for existing checkpoint<br \/>\ncheckpoint_path = \u201cs3:\/\/emr-checkpoints\/job\/progress\/\u201d<br \/>\ns3 = boto3.client(\u2018s3\u2032)<\/p>\n<p>try:<br \/>\ns3.head_object(Bucket=\u2019emr-checkpoints\u2019, Key=\u2019job\/progress\/_SUCCESS\u2019)<br \/>\nprint(\u201cResuming from checkpoint\u2026\u201d)<br \/>\ndf = spark.read.parquet(checkpoint_path)<br \/>\nexcept:<br \/>\nprint(\u201cStarting fresh\u2026\u201d)<br \/>\ndf = spark.read.parquet(\u201cs3:\/\/data-bucket\/raw\/\u201d)<\/p>\n<p># Process with periodic checkpointing<br \/>\ndf = df.filter(df[\u2018date\u2019] &gt;= \u20182024-01-01\u2019)<br \/>\ndf.write.mode(\u2018overwrite\u2019).parquet(checkpoint_path)<\/p>\n<p># Final output<br \/>\ndf.write.mode(\u2018overwrite\u2019).parquet(\u201cs3:\/\/data-bucket\/processed\/\u201d)<\/p>\n<p># Mark completion<br \/>\ns3.put_object(Bucket=\u2019emr-checkpoints\u2019, Key=\u2019job\/progress\/_SUCCESS\u2019, Body=b\u201d)<\/p>\n<p>if __name__ == \u201c__main__\u201d:<br \/>\ntry:<br \/>\nmain()<br \/>\nexcept Exception as e:<br \/>\nprint(f\u201dJob failed: {e}\u201d)<br \/>\nsys.exit(1)<\/p>\n<h3>Cost Monitoring<\/h3>\n<p># Track spot vs on-demand costs per job<br \/>\nemr = boto3.client(\u2019emr\u2019)<br \/>\npricing = boto3.client(\u2018pricing\u2019)<\/p>\n<p>def get_job_cost(cluster_id):<br \/>\n\u201c\u201d\u201dCalculate actual cost for EMR job including spot savings.\u201d\u201d\u201d<\/p>\n<p>instances = emr.list_instances(ClusterId=cluster_id)<br \/>\ntotal_cost = 0<\/p>\n<p>for inst in instances[\u2018Instances\u2019]:<br \/>\ninstance_type = inst[\u2018InstanceType\u2019]<br \/>\nmarket = inst[\u2018Market\u2019] # ON_DEMAND or SPOT<br \/>\nhours = (inst[\u2018Status\u2019][\u2018Timeline\u2019][\u2018EndDateTime\u2019] \u2013<br \/>\ninst[\u2018Status\u2019][\u2018Timeline\u2019][\u2018CreationDateTime\u2019]).seconds \/ 3600<\/p>\n<p># Get pricing (simplified\u2014use CUR for actuals)<br \/>\nif market == \u2018SPOT\u2019:<br \/>\nrate = get_spot_price(instance_type)<br \/>\nelse:<br \/>\nrate = get_ondemand_price(instance_type)<\/p>\n<p>total_cost += rate * hours<\/p>\n<p>return total_cost<\/p>\n<h3>Operational Impact to Measure<\/h3>\n<p>The value of this pattern should be measured with operational metrics rather than treated as an architecture preference. Useful measures include:<\/p>\n<p><strong>When Custom AMIs Don\u2019t Make Sense<\/strong><\/p>\n<ul>\n<li>Ephemeral Clusters (&lt;10 minustes runtime): AMI build time not worth it<\/li>\n<li>Highly Dynamic Dependencies: If requirements change daily<\/li>\n<li>Small Scale (&lt;5 jobs\/day): Bootstrap simplicity wins<\/li>\n<\/ul>\n<h3>Complete Repository Structure<\/h3>\n<p>emr-custom-ami-pipeline\/<br \/>\n\u251c\u2500\u2500 packer\/<br \/>\n\u2502 \u251c\u2500\u2500 emr-custom-ami.pkr.hcl<br \/>\n\u2502 \u251c\u2500\u2500 variables.pkrvars.hcl<br \/>\n\u2502 \u2514\u2500\u2500 scripts\/<br \/>\n\u2502 \u251c\u2500\u2500 install-python.sh<br \/>\n\u2502 \u2514\u2500\u2500 verify-packages.sh<br \/>\n\u251c\u2500\u2500 terraform\/<br \/>\n\u2502 \u251c\u2500\u2500 lambda-builder.tf<br \/>\n\u2502 \u251c\u2500\u2500 step-functions.tf<br \/>\n\u2502 \u251c\u2500\u2500 iam-roles.tf<br \/>\n\u2502 \u2514\u2500\u2500 dynamodb-ami-registry.tf<br \/>\n\u251c\u2500\u2500 lambda\/<br \/>\n\u2502 \u251c\u2500\u2500 ami_builder.py<br \/>\n\u2502 \u251c\u2500\u2500 ami_registry.py<br \/>\n\u2502 \u251c\u2500\u2500 get_golden_ami.py<br \/>\n\u2502 \u251c\u2500\u2500 analyze_emr_failure.py<br \/>\n\u2502 \u2514\u2500\u2500 calculate_backoff.py<br \/>\n\u251c\u2500\u2500 stepfunctions\/<br \/>\n\u2502 \u2514\u2500\u2500 emr-pipeline-asl.json<br \/>\n\u2514\u2500\u2500 spark\/<br \/>\n\u2514\u2500\u2500 process_data.py<\/p>\n<h3>Conclusion<\/h3>\n<p>Custom AMIs can make EMR pipelines more predictable when dependency setup is a recurring source of startup failures. The useful pattern combines:<\/p>\n<ul>\n<li>Immutable dependencies (no PyPI at runtime)<\/li>\n<li>Intelligent retry logic (spot-aware, exponential backoff)<\/li>\n<li>Checkpoint-based recovery (resume from termination)<\/li>\n<\/ul>\n<p>The outcome is not that every EMR workload needs a custom AMI. The outcome is that dependency setup, patching and recovery become explicit platform concerns instead of repeated per-job setup steps.<\/p>\n<p>For teams running frequent Spark pipelines with stable dependency sets, the AMI build investment can pay for itself through fewer environment failures, faster iterations and simpler maintenance.<\/p>\n<p><a href=\"https:\/\/devops.com\/building-reliable-emr-pipelines-with-custom-amis-and-step-functions\/\" target=\"_blank\" class=\"feedzy-rss-link-icon\">Read More<\/a><\/p>\n<p>\u200b<\/p>","protected":false},"excerpt":{"rendered":"<p>A practical pattern for replacing runtime bootstrap fragility with image-based dependencies, orchestration and retry-aware recovery. The Problem: Dependency Drift in [&hellip;]<\/p>\n","protected":false},"author":1,"featured_media":4636,"comment_status":"","ping_status":"","sticky":false,"template":"","format":"standard","meta":{"site-sidebar-layout":"default","site-content-layout":"","ast-site-content-layout":"default","site-content-style":"default","site-sidebar-style":"default","ast-global-header-display":"","ast-banner-title-visibility":"","ast-main-header-display":"","ast-hfb-above-header-display":"","ast-hfb-below-header-display":"","ast-hfb-mobile-header-display":"","site-post-title":"","ast-breadcrumbs-content":"","ast-featured-img":"","footer-sml-layout":"","ast-disable-related-posts":"","theme-transparent-header-meta":"","adv-header-id-meta":"","stick-header-meta":"","header-above-stick-meta":"","header-main-stick-meta":"","header-below-stick-meta":"","astra-migrate-meta-layouts":"default","ast-page-background-enabled":"default","ast-page-background-meta":{"desktop":{"background-color":"var(--ast-global-color-4)","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""},"tablet":{"background-color":"","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""},"mobile":{"background-color":"","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""}},"ast-content-background-meta":{"desktop":{"background-color":"var(--ast-global-color-5)","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""},"tablet":{"background-color":"var(--ast-global-color-5)","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""},"mobile":{"background-color":"var(--ast-global-color-5)","background-image":"","background-repeat":"repeat","background-position":"center center","background-size":"auto","background-attachment":"scroll","background-type":"","background-media":"","overlay-type":"","overlay-color":"","overlay-opacity":"","overlay-gradient":""}},"footnotes":""},"categories":[5],"tags":[],"class_list":["post-4635","post","type-post","status-publish","format-standard","has-post-thumbnail","hentry","category-devops"],"_links":{"self":[{"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/posts\/4635","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/comments?post=4635"}],"version-history":[{"count":0,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/posts\/4635\/revisions"}],"wp:featuredmedia":[{"embeddable":true,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/media\/4636"}],"wp:attachment":[{"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/media?parent=4635"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/categories?post=4635"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/rssfeedtelegrambot.bnaya.co.il\/index.php\/wp-json\/wp\/v2\/tags?post=4635"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}