SageMaker批量转换作业中的竞态条件问题求助
SageMaker Batch Transform任务重复执行最终失败问题排查与解决
我们在生产环境遇到紧急问题:使用SageMaker Batch Transform执行模型推理,每个作业通过ECR中的Docker镜像创建实例,由PyTorch脚本处理任务,完成后调用API存储结果。但CloudWatch日志显示单个任务被重复执行,多次重复后实例无法完成,最终整个操作返回错误。
日志示例
2022-04-24 19:41:47,865 [INFO ] W-model-1-stdout com.amazonaws.ml.mms.wlm.WorkerLifeCycle - Starting to Process Task: 12345678-abcd-1234-efgh-123456ab12c3 ... [任务执行中,输出日志] ... [无错误输出,但任务不再推进,随后同一任务重新启动] ... 2022-04-24 19:52:09,522 [INFO ] W-model-1-stdout com.amazonaws.ml.mms.wlm.WorkerLifeCycle - Starting to Process Task: 12345678-abcd-1234-efgh-123456ab12c3 ... [任务执行中,输出日志] ... [无错误输出,但任务不再推进,随后同一任务重新启动] ... 2022-04-24 20:12:11,834 [INFO ] W-model-1-stdout com.amazonaws.ml.mms.wlm.WorkerLifeCycle - Starting to Process Task: 12345678-abcd-1234-efgh-123456ab12c3 ... [任务执行中,输出日志] ... [无错误输出,但CloudWatch日志停止,Sagemaker向客户端返回错误]
执行作业的代码示例
def inference_batch(self): batch_input = f"s3://{self.cnf.SAGEMAKER_BUCKET}/batch-input/batch.csv" batch_output = f"s3://{self.cnf.SAGEMAKER_BUCKET}/batch-output/" job_name = f"{self.cnf.SAGEMAKER_MODEL}-{str(datetime.datetime.now().strftime('%Y-%m-%d-%H-%m-%S'))}" transform_input = { 'DataSource': { 'S3DataSource': { 'S3DataType': 'S3Prefix', 'S3Uri': batch_input } }, 'ContentType': 'text/csv', 'SplitType': 'Line', } transform_output = { 'S3OutputPath': batch_output } transform_resources = { 'InstanceType': self.cnf.SAGEMAKER_BATCH_INSTANCE, 'InstanceCount': 1 } # self.sm_boto_client 是 boto3.Session(region_name="some-region).client("sagemaker") 的实例 self.sm_boto_client.create_transform_job( TransformJobName=job_name, ModelName=self.cnf.SAGEMAKER_MODEL, TransformInput=transform_input, TransformOutput=transform_output, TransformResources=transform_resources ) status = self.sm_boto_client.describe_transform_job(TransformJobName=job_name) print(f'执行转换任务 {job_name}...') while status['TransformJobStatus'] == 'InProgress': time.sleep(5) status = self.sm_boto_client.describe_transform_job(TransformJobName=job_name) if status['TransformJobStatus'] == 'Completed': print(f'批量转换任务 {job_name} 执行成功。') else: raise Exception(f'批量转换任务 {job_name} 执行失败。')
可能的原因与解决办法
任务超时未触发心跳
SageMaker Batch Transform依赖任务进程的心跳信号判断存活状态,如果PyTorch脚本长时间无输出(比如处理大负载时阻塞过久),Worker生命周期管理会认为进程挂掉,重启任务。- 解决:在PyTorch脚本中增加周期性日志输出(比如每处理N条数据打印一次进度),确保CloudWatch能持续收到心跳信号;调整
MaxConcurrentTransforms和MaxPayloadInMB参数,拆分任务负载,避免单批次处理数据量过大。
- 解决:在PyTorch脚本中增加周期性日志输出(比如每处理N条数据打印一次进度),确保CloudWatch能持续收到心跳信号;调整
资源不足导致进程挂起
实例内存、CPU或GPU资源耗尽,导致PyTorch脚本进程无响应,Worker进程重启任务。- 解决:检查CloudWatch中的系统指标(CPU使用率、内存占用、GPU显存)确认资源瓶颈;升级实例类型或增加
InstanceCount并行处理;在Docker镜像中优化PyTorch脚本的内存使用(比如梯度释放、批量加载数据)。
- 解决:检查CloudWatch中的系统指标(CPU使用率、内存占用、GPU显存)确认资源瓶颈;升级实例类型或增加
MMS配置不合理
SageMaker默认使用MMS管理模型推理进程,若MMS的worker超时配置不合理,会导致任务被重复调度。- 解决:在Docker镜像中修改
config.properties配置文件,调整worker_timeout参数延长超时时间;配置max_request_size参数适配大尺寸输入数据。
- 解决:在Docker镜像中修改
外部API调用阻塞
PyTorch脚本完成推理后调用外部API存储结果,如果API响应缓慢或超时,会导致进程挂起触发Worker重启。- 解决:给API调用添加超时机制(比如
requests库的timeout参数);将结果存储逻辑异步化(比如写入SQS队列由单独服务处理),让推理进程快速完成任务。
- 解决:给API调用添加超时机制(比如
隐性错误触发重试
SageMaker Batch Transform默认有重试逻辑,若任务过程中出现隐性错误(比如S3权限问题、临时网络波动),会触发任务重试。- 解决:通过
describe_transform_job接口查看FailureReason获取具体错误;确保S3输入输出路径的权限配置正确;配置BatchStrategy参数为MultiRecord优化批量处理策略。
- 解决:通过
内容的提问来源于stack exchange,提问作者siavashk
相关产品推荐
相关产品推荐

