You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求使用Python启动AWS EMR集群运行PySpark步骤并自动终止的完整示例

Got it, I’ve been in your shoes before—scouring docs and scattered examples trying to piece together a full, working EMR + PySpark workflow that cleans up after itself. Let me share a complete, tested solution using boto3 (AWS’s Python SDK) that does exactly what you need: spins up an EMR cluster, runs your PySpark job, and terminates automatically whether the job succeeds or fails.


Prerequisites First

Before diving into code, make sure you have these set up:

  • Install boto3: pip install boto3
  • Configure AWS credentials (via ~/.aws/credentials file, environment variables, or IAM roles if running on EC2/EKS)
  • Ensure you have the required IAM roles:
    • Run aws emr create-default-roles to create the default EMR_DefaultRole and EMR_EC2_DefaultRole (these handle cluster permissions)
    • Your user/role needs permissions for EMR (e.g., emr:RunJobFlow, emr:DescribeCluster, emr:ListSteps) and S3 (read/write access for your script, input data, logs, and output)

Full Working Code Example

This script creates an EMR cluster, submits a PySpark step, and ensures the cluster terminates as soon as the step finishes (success or failure):

import boto3
import time

def launch_auto_terminate_emr_cluster():
    # Initialize EMR client (replace region with your preferred one)
    emr_client = boto3.client('emr', region_name='us-east-1')

    # Cluster configuration
    cluster_config = {
        "Name": "Auto-Terminate-PySpark-Workflow",
        "ReleaseLabel": "emr-6.13.0",  # EMR version with Spark 3.4.1 (check AWS docs for compatibility)
        "Instances": {
            "InstanceGroups": [
                {
                    "Name": "Master Node",
                    "InstanceRole": "MASTER",
                    "InstanceType": "m5.xlarge",
                    "InstanceCount": 1,
                    "Market": "ON_DEMAND"
                },
                {
                    "Name": "Core Nodes",
                    "InstanceRole": "CORE",
                    "InstanceType": "m5.xlarge",
                    "InstanceCount": 2,
                    "Market": "ON_DEMAND"
                }
            ],
            "KeepJobFlowAliveWhenNoSteps": False,  # CRITICAL: Terminates cluster when steps finish
            "TerminationProtected": False,  # Allow cluster to terminate on failure
            "Ec2KeyName": "your-ec2-key-pair"  # Optional: Remove if you don't need SSH access
        },
        "Steps": [
            {
                "Name": "Run PySpark Job",
                "ActionOnFailure": "CONTINUE",  # Let cluster terminate even if step fails
                "HadoopJarStep": {
                    "Jar": "command-runner.jar",  # AWS's built-in runner for Spark commands
                    "Args": [
                        "spark-submit",
                        "--deploy-mode", "cluster",
                        "s3://your-bucket/path/to/your/pyspark_script.py",  # Your script in S3
                        # Add your script arguments here, e.g.:
                        # "--input", "s3://your-bucket/input-data/",
                        # "--output", "s3://your-bucket/output-results/"
                    ]
                }
            }
        ],
        "JobFlowRole": "EMR_EC2_DefaultRole",
        "ServiceRole": "EMR_DefaultRole",
        "LogUri": "s3://your-bucket/emr-logs/",  # Optional: Store logs for debugging
        "VisibleToAllUsers": True
    }

    # Create the cluster
    response = emr_client.run_job_flow(**cluster_config)
    cluster_id = response["JobFlowId"]
    print(f"EMR Cluster launched with ID: {cluster_id}")

    # Optional: Poll cluster status until completion (useful for debugging)
    print("Waiting for cluster to complete...")
    while True:
        cluster_status = emr_client.describe_cluster(ClusterId=cluster_id)["Cluster"]["Status"]["State"]
        print(f"Current cluster state: {cluster_status}")
        
        if cluster_status in ["TERMINATED", "TERMINATED_WITH_ERRORS", "FAILED"]:
            break
        time.sleep(30)

    # Check final step status
    steps = emr_client.list_steps(ClusterId=cluster_id)["Steps"]
    for step in steps:
        step_name = step["Name"]
        step_state = step["Status"]["State"]
        print(f"Step '{step_name}' finished with status: {step_state}")

if __name__ == "__main__":
    launch_auto_terminate_emr_cluster()

Key Details Explained
  • KeepJobFlowAliveWhenNoSteps=False: This is the magic setting. It tells EMR to terminate the cluster immediately once all steps have run—regardless of whether they succeeded or failed.
  • ActionOnFailure="CONTINUE": Ensures that if your PySpark step fails, the cluster doesn’t get stuck waiting for manual intervention. It will proceed to terminate right after the failed step.
  • PySpark Script Location: Your .py script must live in S3 (EMR clusters can’t access local files on your machine). Here’s a minimal sample script to test with:
    # s3://your-bucket/path/to/your/pyspark_script.py
    from pyspark.sql import SparkSession
    
    def main():
        spark = SparkSession.builder.appName("TestJob").getOrCreate()
        
        # Sample task: Read CSV, filter, write output
        df = spark.read.csv("s3://your-bucket/input-data/sample.csv", header=True, inferSchema=True)
        filtered_df = df.filter(df["score"] > 80)
        filtered_df.write.mode("overwrite").parquet("s3://your-bucket/output/top-scores/")
        
        spark.stop()
    
    if __name__ == "__main__":
        main()
    

Tips for Production
  • Cost Control: EMR charges by the hour, so auto-termination is critical to avoid unexpected bills.
  • Debugging: Use the LogUri setting to store cluster logs in S3—this is your go-to for troubleshooting failed steps.
  • Instance Types: Adjust instance types/counts based on your job’s resource needs (e.g., use r5 instances for memory-heavy jobs).
  • Version Locking: Pin your EMR release label (e.g., emr-6.13.0) instead of using latest to avoid unexpected breaking changes.

内容的提问来源于stack exchange,提问作者Fred R.

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:35:54