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

在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 模式下,父状态机会阻塞直到子状态机调用 SendTaskSuccess
  • resultPath: "$.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:55:24