如何在集群模式下启动Amazon SageMaker PySparkProcessor?
解决PySparkProcessor集群模式运行问题
你遇到的问题是因为spark.submit.deployMode这个配置在SageMaker PySparkProcessor环境中不适用,SageMaker的Spark集群部署逻辑由框架自身管理,不需要通过这个参数设置。以下是正确的处理方式:
- 移除无效配置:删除
configuration中的spark.submit.deployMode设置,这个参数针对原生Spark集群的提交模式,和SageMaker Processing的运行机制不匹配,添加后反而会导致逻辑冲突。 - 确认集群实例配置:你已经设置
instance_count=2,这会让SageMaker自动启动1个主实例(运行Spark Driver)和1个工作实例(运行Spark Executor),默认就是分布式集群模式。 - 验证集群运行状态:可以在CloudWatch日志中搜索
executor相关内容,查看是否有多个executor注册成功;也可以通过你指定的spark_event_logs_s3_uri路径下的Spark事件日志,确认集群的分布式运行详情。 - 自定义集群资源(可选):如果需要调整executor的资源配置,可在
configuration中添加以下参数:configuration = [{ "Classification": "spark-defaults", "Properties": { "spark.executor.instances": "1", # 数值与instance_count-1对应 "spark.executor.cores": "4", "spark.executor.memory": "8g" } }]
修改后的完整代码示例:
from sagemaker.spark.processing import PySparkProcessor import sagemaker default_bucket = "bucket-xxxxxxxxxxx" sagemaker_session = sagemaker.Session(default_bucket=default_bucket) spark_processor = PySparkProcessor( base_job_name = "sm-spark-vamsi-test", role = vamsi_aws_role, instance_count = 2, sagemaker_session = sagemaker_session, instance_type = 'ml.m5.xlarge', max_runtime_in_seconds = 1200, configuration_location = "s3_location/sagemaker_config_loc" ) pyspark_application_code = "s3_location/test.py" spark_event_logs_s3_uri = "s3_location/logs" submit_py_files = ["s3_location/python_test.zip"] # 配置Spark集群资源参数 configuration = [{ "Classification": "spark-defaults", "Properties": { "spark.executor.instances": "1", "spark.executor.cores": "4", "spark.executor.memory": "8g" }}] # 启动PySpark任务 spark_processor.run( submit_app = pyspark_application_code, spark_event_logs_s3_uri = spark_event_logs_s3_uri, configuration = configuration, submit_py_files = submit_py_files, submit_files = ["s3_location/dev.properties"] )
内容的提问来源于stack exchange,提问作者vamsi
相关产品推荐
相关产品推荐

