如何在Step Functions Map分布式Batch作业中获取作业索引?
问题描述
- 原本在Batch作业脚本中通过环境变量
AWS_BATCH_JOB_ARRAY_INDEX获取数组索引,代码如下:
import os # 获取数组索引 int_idx_array = int(os.environ['AWS_BATCH_JOB_ARRAY_INDEX']) print(f'Array index: {int_idx_array}')
- 这段代码在普通Batch作业中运行正常,但在Step Functions使用Map分布式模式启动Batch作业时,因该环境变量不存在而报错,导致所有作业失败。
- Step Function能正常启动对应数量的Batch作业,但均因上述异常失败,其定义如下:
{ "Comment": "Shared Jobs State Machine", "StartAt": "GetListOfFeatures", "States": { "GetListOfFeatures": { "Type": "Task", "Resource": "arn:aws:lambda:<REGION>:<ACCOUNT>:function:lambda-list-starting-feats-function-ae", "Next": "Tuning" }, "Tuning": { "Type": "Task", "Resource": "arn:aws:states:::batch:submitJob.sync", "Parameters": { "JobQueue": "arn:aws:batch:<REGION>:<ACCOUNT>:job-queue/queue-batch-tuning-ae-image-1", "JobDefinition": "arn:aws:batch:<REGION>:<ACCOUNT>:job-definition/job-def-batch-tuning-ae-image-1:1", "JobName": "job-name-batch-tuning-ae-image-1", "ArrayProperties": { "Size": 10 } }, "Next": "ConcatTuning" }, "ConcatTuning": { "Type": "Task", "Resource": "arn:aws:states:::batch:submitJob.sync", "Parameters": { "JobQueue": "arn:aws:batch:<REGION>:<ACCOUNT>:job-queue/queue-batch-concattuning-ae-image-1", "JobDefinition": "arn:aws:batch:<REGION>:<ACCOUNT>:job-definition/job-def-batch-concattuning-ae-image-1:1", "JobName": "job-name-batch-concattuning-ae-image-1", "ArrayProperties": { "Size": 10 } }, "Next": "Map" }, "Map": { "Type": "Map", "ItemProcessor": { "ProcessorConfig": { "Mode": "DISTRIBUTED", "ExecutionType": "STANDARD" }, "StartAt": "SensitivityAnalysis", "States": { "SensitivityAnalysis": { "Type": "Task", "Resource": "arn:aws:states:::batch:submitJob.sync", "Parameters": { "JobDefinition": "arn:aws:batch:<REGION>:<ACCOUNT>:job-definition/job-def-batch-sensitivity-ae-image-5:1", "JobQueue": "arn:aws:batch:<REGION>:<ACCOUNT>:job-queue/queue-batch-sensitivity-ae-image-5", "JobName": "job-name-batch-sensitivity-ae-image-1" }, "End": true } } }, "End": true, "Label": "Map", "MaxConcurrency": 1000, "ItemReader": { "Resource": "arn:aws:states:::s3:getObject", "ReaderConfig": { "InputType": "CSV", "CSVHeaderLocation": "FIRST_ROW" }, "Parameters": { "Bucket": "20230321-step-functions-poc", "Key": "ae/01_list_of_starting_features/df_cols_in_model.csv" } } } } }
- 当前从S3读取一个CSV文件,行数与需运行的Batch作业数量一致,想了解:
- 如何从该CSV文件中获取索引信息?
- 将CSV保存为键为字符串索引(如"0")的JSON文件,使用该文件是否更简便?
解决方案
核心思路是让Step Functions把当前处理项的索引传递给Batch作业,替代原本依赖的Batch数组环境变量,以下分两种方案说明:
一、使用现有CSV文件的方案
1. 修改Step Function的Map任务配置
在Map状态的SensitivityAnalysis任务中,添加容器环境变量配置,传入Step Functions内置的Map项索引:
"SensitivityAnalysis": { "Type": "Task", "Resource": "arn:aws:states:::batch:submitJob.sync", "Parameters": { "JobDefinition": "arn:aws:batch:<REGION>:<ACCOUNT>:job-definition/job-def-batch-sensitivity-ae-image-5:1", "JobQueue": "arn:aws:batch:<REGION>:<ACCOUNT>:job-queue/queue-batch-sensitivity-ae-image-5", "JobName": "job-name-batch-sensitivity-ae-image-1", "ContainerOverrides": { "Environment": [ { "Name": "ITEM_INDEX", "Value.$": "$$.Map.Item.Index" } ] } }, "End": true }
$$.Map.Item.Index是Step Functions内置变量,代表当前Map项的索引(从0开始递增)。
2. 修改Python代码兼容两种场景
把代码改成优先读取自定义的ITEM_INDEX环境变量,同时保留对原AWS_BATCH_JOB_ARRAY_INDEX的兼容:
import os # 获取索引(兼容普通Batch数组和Step Functions Map模式) int_idx_array = int(os.environ.get('ITEM_INDEX', os.environ.get('AWS_BATCH_JOB_ARRAY_INDEX', 0))) print(f'Array index: {int_idx_array}')
二、使用键为字符串索引的JSON文件的方案
如果已经将CSV转成了键为字符串索引的JSON(如{"0": "feature1", "1": "feature2"}),可以按以下步骤配置:
1. 修改ItemReader读取JSON文件
替换原CSV读取配置为JSON读取:
"ItemReader": { "Resource": "arn:aws:states:::s3:getObject", "ReaderConfig": { "InputType": "JSON" }, "Parameters": { "Bucket": "20230321-step-functions-poc", "Key": "ae/01_list_of_starting_features/your-indexed-file.json" } }
2. 在Batch任务中传递索引和对应值
修改SensitivityAnalysis的参数,同时传入JSON的键(索引)和对应值:
"SensitivityAnalysis": { "Type": "Task", "Resource": "arn:aws:states:::batch:submitJob.sync", "Parameters": { "JobDefinition": "arn:aws:batch:<REGION>:<ACCOUNT>:job-definition/job-def-batch-sensitivity-ae-image-5:1", "JobQueue": "arn:aws:batch:<REGION>:<ACCOUNT>:job-queue/queue-batch-sensitivity-ae-image-5", "JobName": "job-name-batch-sensitivity-ae-image-1", "ContainerOverrides": { "Environment": [ { "Name": "ITEM_INDEX", "Value.$": "$$.Map.Item.Key" }, { "Name": "ITEM_VALUE", "Value.$": "$$.Map.Item.Value" } ] } }, "End": true }
$$.Map.Item.Key对应JSON中的字符串索引,$$.Map.Item.Value对应索引关联的业务数据。
3. Python代码读取索引
使用和CSV方案相同的Python代码即可,无需额外修改。
两种方案对比
- 现有CSV方案更简便,无需转换文件,仅通过Step Functions内置变量就能获取索引,满足基础需求。
- JSON索引文件方案适合需要同时传递索引和关联业务数据的场景,如果仅需索引,没必要额外转换文件。
内容的提问来源于stack exchange,提问作者Aaron England
相关产品推荐
相关产品推荐

