在AWS CDK中获取父Step Function调用的子Step Function输出
AWS CDK 父子 Step Functions 任务令牌集成配置
1. 子 Step Function 发送任务完成信号的配置
子 Step Function 需要在所有业务任务执行完成后,调用 AWS Step Functions 的 SendTaskSuccess API,将任务令牌和执行结果返回给父 Step Function。具体实现如下:
核心配置要点
- 在子状态机末尾添加
CallAwsService任务,调用stepfunctions:SendTaskSuccess接口 - 为子状态机的执行角色添加
states:SendTaskSuccess权限(可限定父状态机 ARN 提升安全性) - 从子状态机输入中获取传入的
taskToken,将子任务执行结果作为output参数传递
CDK 代码示例(子状态机)
import * as cdk from 'aws-cdk-lib'; import * as sfn from 'aws-cdk-lib/aws-stepfunctions'; import * as tasks from 'aws-cdk-lib/aws-stepfunctions-tasks'; import * as iam from 'aws-cdk-lib/aws-iam'; import * as dynamodb from 'aws-cdk-lib/aws-dynamodb'; export class ChildStateMachineStack extends cdk.Stack { public readonly statemachine: sfn.StateMachine; constructor(scope: cdk.App, id: string, props?: cdk.StackProps) { super(scope, id, props); // 子状态机执行角色,添加 SendTaskSuccess 权限 const childRole = new iam.Role(this, 'ChildStateMachineRole', { assumedBy: new iam.ServicePrincipal('states.amazonaws.com'), }); childRole.addToPolicy(new iam.PolicyStatement({ actions: ['states:SendTaskSuccess'], resources: ['*'], // 建议替换为父状态机 ARN 缩小权限范围 })); // 示例业务任务:DynamoDB 写入操作 const targetTable = new dynamodb.Table(this, 'TargetTable', { partitionKey: { name: 'id', type: dynamodb.AttributeType.STRING }, removalPolicy: cdk.RemovalPolicy.DESTROY, }); const dynamoTask = new tasks.DynamoPutItem(this, 'DynamoDB Put Item', { table: targetTable, item: { id: tasks.DynamoAttributeValue.fromString(sfn.JsonPath.stringAt('$.key')), bucket: tasks.DynamoAttributeValue.fromString(sfn.JsonPath.stringAt('$.bucket')), }, resultPath: '$.dynamoResult', }); // 发送任务完成信号的步骤 const sendSuccess = new tasks.CallAwsService(this, 'Send Task Success', { service: 'stepfunctions', action: 'sendTaskSuccess', parameters: { TaskToken: sfn.JsonPath.stringAt('$.token'), Output: sfn.JsonPath.stringAt('$.dynamoResult'), // 传递子任务输出结果 }, iamResources: ['*'], }); // 构建子状态机流程:业务任务 → 发送完成信号 const childDefinition = dynamoTask.next(sendSuccess); this.statemachine = new sfn.StateMachine(this, 'ChildStateMachine', { definitionBody: sfn.DefinitionBody.fromChainable(childDefinition), role: childRole, }); } }
2. 父 Step Function 访问子任务输出的配置
你已在父状态机中设置 resultPath: "$.result",该配置会将子状态机返回的结果存储到父状态机输入的 $.result 路径下,后续状态可直接通过此路径访问输出。
核心说明
IntegrationPattern.WAIT_FOR_TASK_TOKEN模式下,父状态机会阻塞直到子状态机调用SendTaskSuccessresultPath: "$.result"会将子任务结果追加到父输入对象中(而非覆盖整个输入)- 后续状态可通过
sfn.JsonPath.stringAt('$.result')或直接引用$.result获取子任务输出
CDK 代码示例(父状态机后续处理)
import * as cdk from 'aws-cdk-lib'; import * as sfn from 'aws-cdk-lib/aws-stepfunctions'; import * as tasks from 'aws-cdk-lib/aws-stepfunctions-tasks'; import * as iam from 'aws-cdk-lib/aws-iam'; import { ChildStateMachineStack } from './child-stack'; export class ParentStateMachineStack extends cdk.Stack { private executeChild: tasks.StepFunctionsStartExecution; constructor(scope: cdk.App, id: string, props?: cdk.StackProps) { super(scope, id, props); const childStack = new ChildStateMachineStack(this, 'ChildStack'); // 父状态机执行角色,添加启动子状态机的权限 const parentRole = new iam.Role(this, 'ParentStateMachineRole', { assumedBy: new iam.ServicePrincipal('states.amazonaws.com'), }); parentRole.addToPolicy(new iam.PolicyStatement({ actions: ['states:StartExecution'], resources: [childStack.statemachine.stateMachineArn], })); // 调用子状态机的步骤(你的原有代码) this.executeChild = new tasks.StepFunctionsStartExecution(this, "Execute Step Function", { stateMachine: childStack.statemachine, input: sfn.TaskInput.fromObject({ key: sfn.JsonPath.stringAt("$.data.key"), bucket: sfn.JsonPath.stringAt("$.data.bucket"), token: sfn.JsonPath.taskToken, }), resultPath: "$.result", taskTimeout: sfn.Timeout.duration(cdk.Duration.minutes(15)), integrationPattern: sfn.IntegrationPattern.WAIT_FOR_TASK_TOKEN, role: parentRole, }); // 后续处理子任务输出的步骤 const processChildResult = new sfn.Pass(this, 'Process Child Result', { parameters: { childDynamoOutput: sfn.JsonPath.stringAt('$.result'), originalSourceKey: sfn.JsonPath.stringAt('$.data.key'), }, resultPath: '$.processedData', }); // 构建父状态机完整流程 const parentDefinition = this.executeChild.next(processChildResult); new sfn.StateMachine(this, 'ParentStateMachine', { definitionBody: sfn.DefinitionBody.fromChainable(parentDefinition), role: parentRole, }); } }
内容的提问来源于stack exchange,提问作者Luke
相关产品推荐
相关产品推荐

