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

求StepFunction通过Lambda触发EKS多阶段任务的CloudFormation脚本

AWS CloudFormation 模板:StepFunction + Lambda + EKS 多阶段作业流程

以下是实现你需求的CloudFormation模板,包含StepFunction状态机、中间Lambda层、必要IAM权限及EKS作业执行配置,可直接替换占位符后部署:

AWSTemplateFormatVersion: '2010-09-09'
Description: 部署StepFunction+Lambda+EKS多阶段作业流程

Parameters:
  EKSClusterName:
    Type: String
    Description: 目标EKS集群名称
  ECRImage1URI:
    Type: String
    Description: Stage1对应的ECR镜像URI
  ECRImage2URI:
    Type: String
    Description: Stage2对应的ECR镜像URI
  ECRImage3URI:
    Type: String
    Description: Stage3对应的ECR镜像URI
  EKSNamespace:
    Type: String
    Default: default
    Description: EKS中运行作业的命名空间

Resources:
  # ------------------------------
  # IAM角色:StepFunction调用Lambda权限
  # ------------------------------
  StepFunctionExecutionRole:
    Type: AWS::IAM::Role
    Properties:
      AssumeRolePolicyDocument:
        Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Principal:
              Service: states.amazonaws.com
            Action: sts:AssumeRole
      Policies:
        - PolicyName: InvokeLambdaPolicy
          PolicyDocument:
            Version: '2012-10-17'
            Statement:
              - Effect: Allow
                Action: lambda:InvokeFunction
                Resource: !Ref EKSJobLambdaFunction

  # ------------------------------
  # IAM角色:Lambda调用EKS及相关权限
  # ------------------------------
  LambdaExecutionRole:
    Type: AWS::IAM::Role
    Properties:
      AssumeRolePolicyDocument:
        Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Principal:
              Service: lambda.amazonaws.com
            Action: sts:AssumeRole
      ManagedPolicyArns:
        - arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole
      Policies:
        - PolicyName: EKSClusterAccessPolicy
          PolicyDocument:
            Version: '2012-10-17'
            Statement:
              - Effect: Allow
                Action: eks:DescribeCluster
                Resource: !Sub arn:aws:eks:${AWS::Region}:${AWS::AccountId}:cluster/${EKSClusterName}
        - PolicyName: EKSJobExecutionPolicy
          PolicyDocument:
            Version: '2012-10-17'
            Statement:
              - Effect: Allow
                Action:
                  - eks:DescribeFargateProfile
                  - eks:ListFargateProfiles
                Resource: !Sub arn:aws:eks:${AWS::Region}:${AWS::AccountId}:fargateprofile/${EKSClusterName}/*

  # ------------------------------
  # Lambda函数:中间层调用EKS创建并等待Job完成
  # ------------------------------
  EKSJobLambdaFunction:
    Type: AWS::Lambda::Function
    Properties:
      Handler: index.lambda_handler
      Runtime: python3.11
      Role: !GetAtt LambdaExecutionRole.Arn
      Code:
        ZipFile: |
          import boto3
          import kubernetes
          from kubernetes.client.rest import ApiException
          import time
          import os

          def lambda_handler(event, context):
              # 初始化EKS客户端获取kubeconfig
              eks_client = boto3.client('eks')
              cluster_name = os.environ['EKS_CLUSTER_NAME']
              namespace = os.environ['EKS_NAMESPACE']
              job_name = event['JobName']
              image_uri = event['ImageURI']

              try:
                  cluster_info = eks_client.describe_cluster(name=cluster_name)
                  kubeconfig = {
                      'apiVersion': 'v1',
                      'clusters': [{
                          'cluster': {
                              'server': cluster_info['cluster']['endpoint'],
                              'certificate-authority-data': cluster_info['cluster']['certificateAuthority']['data']
                          },
                          'name': cluster_name
                      }],
                      'contexts': [{
                          'context': {
                              'cluster': cluster_name,
                              'user': f'aws-lambda-{cluster_name}'
                          },
                          'name': cluster_name
                      }],
                      'current-context': cluster_name,
                      'kind': 'Config',
                      'users': [{
                          'name': f'aws-lambda-{cluster_name}',
                          'user': {
                              'exec': {
                                  'apiVersion': 'client.authentication.k8s.io/v1beta1',
                                  'command': 'aws-iam-authenticator',
                                  'args': [
                                      'token',
                                      '-i', cluster_name
                                  ]
                              }
                          }
                      }]
                  }

                  # 配置K8s客户端
                  kubernetes.config.load_kube_config_from_dict(kubeconfig)
                  api_instance = kubernetes.client.BatchV1Api()

                  # 定义Job模板
                  job_body = kubernetes.client.V1Job(
                      api_version="batch/v1",
                      kind="Job",
                      metadata=kubernetes.client.V1ObjectMeta(name=job_name),
                      spec=kubernetes.client.V1JobSpec(
                          template=kubernetes.client.V1PodTemplateSpec(
                              spec=kubernetes.client.V1PodSpec(
                                  restart_policy="Never",
                                  containers=[
                                      kubernetes.client.V1Container(
                                          name=job_name,
                                          image=image_uri
                                      )
                                  ]
                              )
                          ),
                          backoff_limit=0
                      )
                  )

                  # 创建Job
                  api_instance.create_namespaced_job(namespace=namespace, body=job_body)
                  print(f"Job {job_name} 创建成功")

                  # 轮询等待Job完成
                  while True:
                      time.sleep(10)
                      job_status = api_instance.read_namespaced_job_status(name=job_name, namespace=namespace)
                      if job_status.status.succeeded:
                          print(f"Job {job_name} 执行完成")
                          return {'Status': 'SUCCEEDED', 'JobName': job_name}
                      elif job_status.status.failed:
                          raise Exception(f"Job {job_name} 执行失败")

              except ApiException as e:
                  raise Exception(f"K8s API调用失败: {e}")
              except Exception as e:
                  raise Exception(f"Lambda执行失败: {str(e)}")
      Environment:
        Variables:
          EKS_CLUSTER_NAME: !Ref EKSClusterName
          EKS_NAMESPACE: !Ref EKSNamespace
      Timeout: 900 # 15分钟超时,可根据作业时长调整

  # ------------------------------
  # StepFunction状态机:编排三个阶段的作业流程
  # ------------------------------
  EKSJobStateMachine:
    Type: AWS::StepFunctions::StateMachine
    Properties:
      RoleArn: !GetAtt StepFunctionExecutionRole.Arn
      DefinitionString: !Sub |
        {
          "Comment": "多阶段EKS作业执行流程",
          "StartAt": "Stage1",
          "States": {
            "Stage1": {
              "Type": "Task",
              "Resource": "${EKSJobLambdaFunction.Arn}",
              "Parameters": {
                "JobName": "stage1-job",
                "ImageURI": "${ECRImage1URI}"
              },
              "Next": "Stage2"
            },
            "Stage2": {
              "Type": "Task",
              "Resource": "${EKSJobLambdaFunction.Arn}",
              "Parameters": {
                "JobName": "stage2-job",
                "ImageURI": "${ECRImage2URI}"
              },
              "Next": "Stage3"
            },
            "Stage3": {
              "Type": "Task",
              "Resource": "${EKSJobLambdaFunction.Arn}",
              "Parameters": {
                "JobName": "stage3-job",
                "ImageURI": "${ECRImage3URI}"
              },
              "End": true
            }
          }
        }

  # ------------------------------
  # EKS服务账号权限:允许拉取ECR镜像
  # ------------------------------
  EKSJobServiceAccount:
    Type: AWS::IAM::Role
    Properties:
      AssumeRolePolicyDocument:
        Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Principal:
              Federated: !Sub arn:aws:iam::${AWS::AccountId}:oidc-provider/${EKSClusterOIDCProvider}
            Action: sts:AssumeRoleWithWebIdentity
            Condition:
              StringEquals:
                !Sub "${EKSClusterOIDCProvider}:sub": !Sub "system:serviceaccount:${EKSNamespace}:eks-job-sa"
      Policies:
        - PolicyName: ECRPullPolicy
          PolicyDocument:
            Version: '2012-10-17'
            Statement:
              - Effect: Allow
                Action:
                  - ecr:GetDownloadUrlForLayer
                  - ecr:BatchGetImage
                  - ecr:BatchCheckLayerAvailability
                  - ecr:GetAuthorizationToken
                Resource: "*"

  # 获取EKS集群OIDC提供商URL
  EKSClusterOIDCProvider:
    Type: AWS::SSM::Parameter::Value<String>
    Default: !Sub /aws/eks/${EKSClusterName}/oidc-issuer

Outputs:
  StateMachineARN:
    Description: StepFunction状态机ARN
    Value: !Ref EKSJobStateMachine
  LambdaFunctionARN:
    Description: Lambda函数ARN
    Value: !Ref EKSJobLambdaFunction

关键配置说明

  1. Lambda函数逻辑

    • 使用boto3获取EKS集群的端点和证书,动态生成kubeconfig
    • 通过kubernetes客户端创建Job,并轮询检查Job状态,直到成功/失败
    • 需确保Lambda运行环境已安装kubernetes和aws-iam-authenticator(可通过Lambda层添加依赖,或使用预打包自定义镜像)
  2. 权限配置要点

    • Lambda角色需具备eks:DescribeCluster权限以获取集群信息
    • 通过EKS OIDC身份提供商,将IAM角色绑定到EKS服务账号,确保Pod能拉取ECR镜像
    • StepFunction角色仅保留调用目标Lambda的权限,遵循最小权限原则
  3. 自定义调整建议

    • 根据作业实际时长修改Lambda的Timeout参数
    • 若需更细粒度的权限,可限制ECR镜像的Resource为具体的ECR仓库ARN
    • 可在StepFunction中添加Catch字段实现错误重试或告警分支

内容的提问来源于stack exchange,提问作者Ravi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 09:47:36