求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
关键配置说明
Lambda函数逻辑
- 使用
boto3获取EKS集群的端点和证书,动态生成kubeconfig - 通过
kubernetes客户端创建Job,并轮询检查Job状态,直到成功/失败 - 需确保Lambda运行环境已安装
kubernetes和aws-iam-authenticator(可通过Lambda层添加依赖,或使用预打包自定义镜像)
- 使用
权限配置要点
- Lambda角色需具备
eks:DescribeCluster权限以获取集群信息 - 通过EKS OIDC身份提供商,将IAM角色绑定到EKS服务账号,确保Pod能拉取ECR镜像
- StepFunction角色仅保留调用目标Lambda的权限,遵循最小权限原则
- Lambda角色需具备
自定义调整建议
- 根据作业实际时长修改Lambda的
Timeout参数 - 若需更细粒度的权限,可限制ECR镜像的Resource为具体的ECR仓库ARN
- 可在StepFunction中添加
Catch字段实现错误重试或告警分支
- 根据作业实际时长修改Lambda的
内容的提问来源于stack exchange,提问作者Ravi
相关产品推荐
相关产品推荐

