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

如何在AWS CDK Step Functions并行任务间访问值及使用Context Object

在AWS CDK的Step Functions并行任务中跨分支访问数据的方案

首先明确:并行分支的任务是同时执行的,在任务运行过程中,一个分支无法直接访问另一个分支的实时数据——Context Object也做不到这一点,它只能访问当前任务/执行流的上下文(比如自身的任务ID、启动时间),无法跨并行分支获取其他任务的信息。

以下是两种可行的解决方案,根据你的业务需求选择:

方案1:并行任务完成后统一处理输出

如果不需要在任务执行过程中互相引用数据,只是需要在并行任务全部结束后使用彼此的结果(比如Glue Job Run ID),可以将两个并行任务的输出整合,在后续步骤中访问。

CDK代码示例

import * as sfn from 'aws-cdk-lib/aws-stepfunctions';
import * as tasks from 'aws-cdk-lib/aws-stepfunctions-tasks';
import * as glue from 'aws-cdk-lib/aws-glue';
import { Stack } from 'aws-cdk-lib';
import { Construct } from 'constructs';

export class ParallelGlueStack extends Stack {
  constructor(scope: Construct, id: string) {
    super(scope, id);

    // 定义两个Glue Job
    const glueJob1 = new glue.Job(this, 'GlueJob1', {
      jobName: 'DataProcessingJob1',
      executable: glue.JobExecutable.pythonEtl({
        glueVersion: glue.GlueVersion.V4_0,
        pythonVersion: glue.PythonVersion.THREE_NINE,
        script: glue.Code.fromAsset('glue-scripts/job1.py'),
      }),
    });

    const glueJob2 = new glue.Job(this, 'GlueJob2', {
      jobName: 'DataProcessingJob2',
      executable: glue.JobExecutable.pythonEtl({
        glueVersion: glue.GlueVersion.V4_0,
        pythonVersion: glue.PythonVersion.THREE_NINE,
        script: glue.Code.fromAsset('glue-scripts/job2.py'),
      }),
    });

    // 创建Glue任务执行步骤,指定结果存储路径
    const runJob1 = new tasks.GlueStartJobRun(this, 'RunJob1', {
      glueJob: glueJob1,
      resultPath: '$.job1Output',
    });

    const runJob2 = new tasks.GlueStartJobRun(this, 'RunJob2', {
      glueJob: glueJob2,
      resultPath: '$.job2Output',
    });

    // 并行执行两个任务,整合输出
    const parallelJobs = new sfn.Parallel(this, 'ParallelGlueJobs', {
      resultPath: '$.parallelResults',
    }).branch(runJob1).branch(runJob2);

    // 后续步骤:访问两个任务的Run ID
    const processResults = new sfn.Pass(this, 'ProcessResults', {
      parameters: {
        Job1RunId: '$.parallelResults[0].RunId',
        Job2RunId: '$.parallelResults[1].RunId',
        // 若需要将其中一个ID传给其他任务,可在此构造参数后传递
      },
    });

    // 组装状态机流程
    const definition = parallelJobs.next(processResults);
    new sfn.StateMachine(this, 'ParallelStateMachine', {
      definitionBody: sfn.DefinitionBody.fromChainable(definition),
      stateMachineName: 'ParallelGlueProcessing',
    });
  }
}

方案2:串行执行任务,提前传递依赖数据

如果必须在某个Glue任务执行时就用到另一个任务的ID(比如Job2启动时需要Job1的Run ID作为参数),则无法使用并行结构,必须调整为串行执行:先启动第一个任务,获取其输出后,再启动第二个任务并传入所需参数。

CDK代码示例

import * as sfn from 'aws-cdk-lib/aws-stepfunctions';
import * as tasks from 'aws-cdk-lib/aws-stepfunctions-tasks';
import * as glue from 'aws-cdk-lib/aws-glue';
import { Stack } from 'aws-cdk-lib';
import { Construct } from 'constructs';

export class SerialGlueStack extends Stack {
  constructor(scope: Construct, id: string) {
    super(scope, id);

    const glueJob1 = new glue.Job(this, 'GlueJob1', {
      jobName: 'DataProcessingJob1',
      executable: glue.JobExecutable.pythonEtl({
        glueVersion: glue.GlueVersion.V4_0,
        pythonVersion: glue.PythonVersion.THREE_NINE,
        script: glue.Code.fromAsset('glue-scripts/job1.py'),
      }),
    });

    const glueJob2 = new glue.Job(this, 'GlueJob2', {
      jobName: 'DataProcessingJob2',
      executable: glue.JobExecutable.pythonEtl({
        glueVersion: glue.GlueVersion.V4_0,
        pythonVersion: glue.PythonVersion.THREE_NINE,
        script: glue.Code.fromAsset('glue-scripts/job2.py'),
      }),
    });

    // 先执行Job1,获取Run ID
    const runJob1 = new tasks.GlueStartJobRun(this, 'RunJob1', {
      glueJob: glueJob1,
      resultPath: '$.job1Output',
    });

    // 将Job1的Run ID作为参数传入Job2
    const runJob2 = new tasks.GlueStartJobRun(this, 'RunJob2', {
      glueJob: glueJob2,
      arguments: sfn.TaskInput.fromObject({
        '--JOB1_RUN_ID': '$.job1Output.RunId',
      }),
      resultPath: '$.job2Output',
    });

    // 组装串行流程
    const definition = runJob1.next(runJob2);
    new sfn.StateMachine(this, 'SerialStateMachine', {
      definitionBody: sfn.DefinitionBody.fromChainable(definition),
      stateMachineName: 'SerialGlueProcessing',
    });
  }
}

关键结论

  • Context Object仅能访问当前任务的上下文数据,无法跨并行分支获取其他任务的信息,不是解决这个问题的正确方案。
  • 严格并行的任务只能在全部完成后,在后续步骤中互相访问输出结果。
  • 若任务间存在实时数据依赖,必须调整为串行执行,先获取依赖数据再启动后续任务。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 04:45:45