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

如何在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作业数量一致,想了解:
    1. 如何从该CSV文件中获取索引信息?
    2. 将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 12:37:08