求使用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/credentialsfile, environment variables, or IAM roles if running on EC2/EKS) - Ensure you have the required IAM roles:
- Run
aws emr create-default-rolesto create the defaultEMR_DefaultRoleandEMR_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)
- Run
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
.pyscript 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
LogUrisetting 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
r5instances for memory-heavy jobs). - Version Locking: Pin your EMR release label (e.g.,
emr-6.13.0) instead of usinglatestto avoid unexpected breaking changes.
内容的提问来源于stack exchange,提问作者Fred R.
相关产品推荐
相关产品推荐

