- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.
导读
本文基于开源仓库 aws-doc-sdk-examples 中kotlin/services/stepfunctions目录的官方示例,系统讲解如何用 AWS SDK for Kotlin 调用 AWS Step Functions 的完整操作集:从列出状态机(Hello 入门)、创建/删除活动与状态机、描述执行、获取执行历史、筛选失败执行,到结合 IAM 角色、活动任务轮询与sendTaskSuccess回调的端到端"聊天模拟器"场景。读完本文,你将掌握 Step Functions 在 Kotlin 中的客户端初始化方式、分页器(Paginator)用法、状态机定义文件的动态注入技巧,以及 JUnit 5 自动化测试的配置方法,可直接运行仓库中的全部示例。
Step Functions 与 SDK for Kotlin 示例概览
AWS Step Functions 是一个可视化工作流服务,帮助开发者使用 AWS 服务构建分布式应用、自动化流程、编排微服务,并创建数据处理和机器学习(ML)管道。本仓库的 Kotlin 示例正是围绕这一核心能力展开,官方 README(kotlin/services/stepfunctions/README.md)将其组织为三个层次:
- Get started(入门):
HelloStepFunctions.kt,演示listStateMachines命令; - Single action(单动作):每个代码片段只调用一个服务函数,覆盖
createActivity、createStateMachine、deleteActivity、deleteStateMachine、describeExecution、describeStateMachine、getActivityTask、getExecutionHistory、listExecutions、sendTaskSuccess、startExecution共 11 个 API; - Scenarios(场景):
StepFunctionsScenario.kt,通过在同一服务内调用多个函数完成完整任务。
所有示例统一使用SfnClient,它是 SDK for Kotlin 提供的 Step Functions 服务客户端。源码中可以看到其标准初始化模式(如 HelloStepFunctions.kt):
SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> // 调用 Step Functions API }fromEnvironment表示从环境变量(AWS_ACCESS_KEY_ID、AWS_SECRET_ACCESS_KEY、AWS_REGION等)加载凭证与配置,use块保证客户端在使用后自动释放底层资源。示例默认使用us-east-1区域,运行前请结合 AWS Regional Services 确认该服务在你所在区域可用。
Hello 入门:列出状态机
HelloStepFunctions.kt 是最简示例,核心逻辑如下:
suspend fun listMachines() { SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> val response = sfnClient.listStateMachines(ListStateMachinesRequest {}) response.stateMachines?.forEach { machine -> println("The name of the state machine is ${machine.name}") println("The ARN value is ${machine.stateMachineArn}") } } }该示例无命令行参数,直接运行即可列出账户中最多十个状态机的名称与 ARN。注意函数是suspend的,SDK for Kotlin 的 API 基于 Kotlin 协程实现,调用方必须在协程作用域内执行(仓库测试中即用runBlocking包裹)。
单动作示例源码剖析
列出活动:ListActivities.kt
ListActivities.kt 使用listActivities接口并设置maxResults = 10:
suspend fun listAllActivites() { val activitiesRequest = ListActivitiesRequest { maxResults = 10 } SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> val response = sfnClient.listActivities(activitiesRequest) response.activities?.forEach { item -> println("The activity ARN is ${item.activityArn}") println("The activity name is ${item.name}") } } }获取执行历史:GetExecutionHistory.kt
GetExecutionHistory.kt 接收一个执行 ARN 作为命令行参数,调用getExecutionHistory拉取执行事件并打印事件类型:
suspend fun getExeHistory(exeARN: String?) { val historyRequest = GetExecutionHistoryRequest { executionArn = exeARN maxResults = 10 } SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> val response = sfnClient.getExecutionHistory(historyRequest) response.events?.forEach { event -> println("The event type is ${event.type}") } } }运行方式:GetExecutionHistoryKt <exeARN>,其中exeARN是某次执行的 ARN(可从startExecution的返回值获得)。
筛选失败执行:GetFailedExecutions.kt
GetFailedExecutions.kt 演示listExecutions与状态过滤的组合,通过statusFilter = ExecutionStatus.Failed只返回失败的执行记录:
suspend fun getFailedExes(stateMachineARN: String?) { val executionsRequest = ListExecutionsRequest { maxResults = 10 stateMachineArn = stateMachineARN statusFilter = ExecutionStatus.Failed } SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> val response = sfnClient.listExecutions(executionsRequest) response.executions?.forEach { item -> println("The Amazon Resource Name (ARN) of the failed execution is ${item.executionArn}.") } } }运行方式:GetFailedExecutionsKt <stateMachineARN>,stateMachineARN即状态机的 ARN。这个示例很适合作为"工作流失败监控"的起点,statusFilter还可替换为Running、Succeeded、TimedOut、Aborted等枚举值。
核心场景:StepFunctionsScenario 端到端实战
StepFunctionsScenario.kt 是仓库中最完整的示例,源码注释明确列出 9 个步骤,覆盖"创建 → 交互 → 清理"全生命周期:
- 使用分页器(Paginator)列出活动;
- 使用分页器列出状态机;
- 创建一个活动(Activity);
- 创建一个状态机;
- 描述该状态机;
- 启动状态机执行并与之交互;
- 描述执行状态;
- 删除活动;
- 删除状态机。
运行参数
该场景运行时需要 4 个命令行参数:
| 参数 | 含义 |
|---|---|
roleName | 为状态机创建的 IAM 角色名称 |
activityName | 要创建的活动名称 |
stateMachineName | 要创建的状态机名称 |
jsonFile | chat_sfn_state_machine.json文件的本地路径 |
动态注入活动 ARN 到状态机定义
场景的核心技巧是用 Jackson 读取 chat_sfn_state_machine.json,再把GetInput状态中的Resource字段替换为刚创建的活动的 ARN,从而让状态机把"等待人类输入"这一任务回调给该活动:
// 读取 JSON 定义文件 val stream = GetStream() val jsonString = stream.getStream(jsonFile) // 将 Resource 节点替换为 activityArn val objectMapper = ObjectMapper() val root: JsonNode = objectMapper.readTree(jsonString) (root.path("States").path("GetInput") as ObjectNode).put("Resource", activityArn) val stateDefinition = objectMapper.writeValueAsString(root)GetStream.kt 负责把本地 JSON 文件读入并序列化为字符串,供后续createStateMachine使用。仓库模板文件 chat_sfn_state_machine.json 中的GetInput状态原本的Resource是占位符{{DOC_EXAMPLE_ACTIVITY_ARN}},其TimeoutSeconds设为 300,表示活动必须在 300 秒内回调,否则任务超时。
创建 IAM 角色
状态机需要执行角色(execution role),场景代码用IamClient创建,信任策略允许states.amazonaws.com服务代入该角色(即sts:AssumeRole):
val polJSON = """{ "Version": "2012-10-17", "Statement": [ { "Sid": "", "Effect": "Allow", "Principal": { "Service": "states.amazonaws.com" }, "Action": "sts:AssumeRole" } ] }""" suspend fun createIAMRole(roleNameVal: String?, polJSON: String?): String? { val request = CreateRoleRequest { roleName = roleNameVal assumeRolePolicyDocument = polJSON description = "Created using the AWS SDK for Kotlin" } IamClient.fromEnvironment { region = "AWS_GLOBAL" }.use { iamClient -> val response = iamClient.createRole(request) return response.role?.arn } }注意 IAM 是全局服务,此处区域使用AWS_GLOBAL,与 Step Functions 的us-east-1不同。务必按最小权限原则:示例仅创建角色本身,生产环境还应为其附加恰好满足状态机所需操作的权限策略。
创建活动与状态机
// 创建活动 suspend fun createActivity(activityName: String): String? { val activityRequest = CreateActivityRequest { name = activityName } SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> val response = sfnClient.createActivity(activityRequest) return response.activityArn } } // 创建标准类型状态机 suspend fun createMachine(roleARNVal: String?, stateMachineName: String?, jsonVal: String?): String? { val machineRequest = CreateStateMachineRequest { definition = jsonVal name = stateMachineName roleArn = roleARNVal type = StateMachineType.Standard } ... }type = StateMachineType.Standard表示标准工作流(适合有状态、需人工介入或长时间运行的流程),如需高吞吐、短时执行可改用StateMachineType.Express。
启动执行与活动任务轮询
这是整个场景最核心的交互部分:先用startExecution携带 JSON 输入启动工作流,再循环调用getActivityTask轮询活动任务,读取状态机下发的消息后由用户在控制台输入回复,最后用sendTaskSuccess把用户输入回传给 Step Functions:
val runArn = startWorkflow(stateMachineArn, executionJson) while (!action) { myList = getActivityTask(activityArn) // 1. 获取活动任务(token + 输入) println("ChatSFN: " + myList[1]) // 2. 打印状态机发来的消息 val myAction = sc.nextLine() // 3. 等待用户输入 if (myAction.compareTo("done") == 0) action = true val taskJson = "{ \"action\" : \"$myAction\" }" sendTaskSuccess(myList[0], taskJson) // 4. 回调 SendTaskSuccess }其中startWorkflow使用UUID.randomUUID()生成执行名,保证并发执行不冲突;getActivityTask返回List<String>,元素 0 是taskToken(后续回调必须携带的令牌),元素 1 是活动输入 JSON。
对照 chat_sfn_state_machine.json 的定义可理解完整闭环:Greeting状态用States.Format拼出问候语 →GetInput任务等待活动回调 →Choice根据$.action分发到 Song/Poem/Story/Done 四个分支 → 非 Done 分支再次回到GetInput,形成对话循环;用户输入done后进入结束状态。这套"Activity + SendTaskSuccess"模式是 Step Functions 实现人工审批、外部系统回调的标准做法。
描述与轮询执行状态
describeStateMachine打印状态机名称、状态、ARN 与角色 ARN;describeExe则轮询describeExecution,每 2 秒检查一次,直到执行状态从Running变为Succeeded:
while (!hasSucceeded) { val response = sfnClient.describeExecution(executionRequest) status = response.status.toString() if (status.compareTo("Running") == 0) { println("The state machine is still running, let's wait for it to finish.") Thread.sleep(2000) } else if (status.compareTo("Succeeded") == 0) { println("The Step Function workflow has succeeded") hasSucceeded = true } else { println("The Status is $status") } }清理资源
场景结尾依次调用deleteActivity(activityArn)与deleteMachine(stateMachineArn)删除活动与状态机,避免遗留资源持续计费。README 特别提醒:运行删除/修改 AWS 资源的操作时要格外小心,建议使用独立的测试专用资源进行实验。
分页器(Paginator)的协程流用法
场景前两步展示了 SDK for Kotlin 的分页器能力——通过listActivitiesPaginated/listStateMachinesPaginated配合kotlinx.coroutines.flow.transform自动翻页并流式处理结果:
suspend fun listStatemachinesPagnator() { val machineRequest = ListStateMachinesRequest { maxResults = 10 } SfnClient.fromEnvironment { region = "us-east-1" }.use { sfnClient -> sfnClient.listStateMachinesPaginated(machineRequest) .transform { it.stateMachines?.forEach { obj -> emit(obj) } } .collect { obj -> println(" The state machine ARN is ${obj.stateMachineArn}") } } }相比手动循环调用listStateMachines并处理nextToken,分页器把分页逻辑封装进Flow,代码更简洁且天然适配 Kotlin 协程。这与 HelloStepFunctions.kt 中的非分页版本形成了"基础 API 与高级 API"的对照学习材料。
构建配置与运行环境
示例项目的构建配置见 build.gradle.kts,要点如下:
- Kotlin JVM 插件版本 2.1.0,Java 兼容级别 17(
sourceCompatibility/targetCompatibility均为VERSION_17,jvmTarget = "17"); - 通过BOM(
aws.sdk.kotlin:bom:1.5.63)统一管理 AWS SDK for Kotlin 版本,无需为各服务模块单独指定版本; - 依赖模块包括
sfn(Step Functions)、iam、secretsmanager(测试读取密钥)、OkHttp/CRT HTTP 引擎、Jackson(JSON 解析)、Gson、kotlinx-coroutines-core、SLF4J 日志; - 测试框架为JUnit Jupiter 5.9.2,并集成
ktlint-gradle插件做代码风格检查; - 测试任务通过
useJUnitPlatform()启用 JUnit Platform。
运行前需先按 AWS SDK for Kotlin 开发者指南 完成开发环境与凭证配置(AWS_ACCESS_KEY_ID、AWS_SECRET_ACCESS_KEY、AWS_REGION或~/.aws/credentials)。运行这些代码可能对您的 AWS 账户产生费用,详见 AWS Pricing。
JUnit 5 自动化测试:StepFunctionsKotlinTest
仓库提供测试文件 StepFunctionsKotlinTest.kt,基于 JUnit 5 编写,可从 IntelliJ 等 IDE 或命令行运行。每跑完一个用例会输出类似Test 1 passed的信息;⚠️ 运行这些测试会操作真实 AWS 资源,可能产生账户费用。
配置属性
README 说明:运行测试前必须在resources文件夹下的config.properties文件中定义以下值,否则测试失败:
| 属性 | 含义 |
|---|---|
roleNameSc | 场景测试要创建的 IAM 角色名称 |
activityNameSc | 场景测试要创建的活动名称 |
stateMachineNameSc | 场景测试要创建的状态机名称 |
仓库中的 config.properties 实际列出的键更完整,包括jsonFile、jsonFileSM、roleARN、stateMachineName以及上述三个*Sc键,均需按实际环境填入。
状态机定义文件
测试状态机使用 chat_sfn_state_machine.json。按 README 要求,需将该项目文件放入你工程(即kotlin/services/stepfunctions)的resources文件夹中;若缺失,场景测试将失败。
当前测试实现:从 Secrets Manager 读取
值得注意,仓库当前的测试实现(StepFunctionsKotlinTest.kt)已升级为从 AWS Secrets Manager 读取配置:@BeforeAll阶段从名为test/stepfunctions的密钥中取 JSON,解析出roleNameSC、activityNameSC、stateMachineNameSC与machineFile四个字段,并在角色名、活动名、状态机名后追加UUID.randomUUID()以保证测试隔离:
@BeforeAll fun setup() = runBlocking { val gson = Gson() val json: String = getSecretValues() // 从 "test/stepfunctions" 密钥读取 val values = gson.fromJson(json, SecretValues::class.java) roleNameSC = values.roleNameSC.toString() + UUID.randomUUID() activityNameSC = values.activityNameSC.toString() + UUID.randomUUID() stateMachineNameSC = values.stateMachineNameSC.toString() + UUID.randomUUID() jsonFile = values.machineFile.toString() }getSecretValues()通过SecretsManagerClient { region = "us-east-1" }调用getSecretValue读取密钥内容。因此,实际跑测试前需要预先在 Secrets Manager 中创建test/stepfunctions密钥并填入上述字段(对应仓库中 build.gradle.kts 对secretsmanager模块的依赖)。测试用例 1 直接调用场景中的listMachines()函数验证 Hello 逻辑,这是源码层面 README 与实现之间最直观的印证。
总结与使用建议
将 README 与仓库源码对照可得出以下实践要点:
- 从 Hello 到场景循序渐进:先用 HelloStepFunctions.kt 验证凭证与连通性,再通过 ListActivities.kt、GetExecutionHistory.kt、GetFailedExecutions.kt 熟悉单 API,最后深入 StepFunctionsScenario.kt 掌握活动回调模式;
- 善用分页器与协程:列表类 API 优先使用
*Paginated变体,配合Flow的transform/collect处理大数据量结果; - 状态机定义可编程注入:用 Jackson 在运行时替换
Resource、参数等节点,是构建通用工作流引擎的常见手段; - 权限最小化:IAM 角色遵循 least privilege 原则,只为状态机授予完成任务所需的最小权限集;清理示例务必执行删除步骤;
- 测试配置二选一:按 README 用
config.properties,或按当前实现预置 Secrets Manager 密钥test/stepfunctions,并在测试环境使用独立资源以避免污染生产账户。
所有源码均位于kotlin/services/stepfunctions目录下,示例版权归 Amazon.com, Inc. 或其关联公司所有,采用 Apache-2.0 许可证(SPDX-License-Identifier: Apache-2.0),可放心学习与复用。
- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.
相关推荐
使用 AWS SDK for Kotlin 操作 AWS IoT Core:场景化示例代码全解析
使用 AWS SDK for Kotlin 操作 AWS IoT Core:场景化示例代码全解析 导读 本文基于 AWS 官方代码示例仓库中 kotlin/se
示例工程教程后端使用 AWS SDK for Kotlin 操作 Amazon EventBridge:从基础示例到完整场景实战
使用 AWS SDK for Kotlin 操作 Amazon EventBridge:从基础示例到完整场景实战 Amazon EventBridge 是 AW
示例工程教程后端使用 AWS SDK for Kotlin 操作 Amazon EC2 Auto Scaling:完整场景示例与测试指南
使用 AWS SDK for Kotlin 操作 Amazon EC2 Auto Scaling:完整场景示例与测试指南 导读 本文基于 aws doc sdk
示例工程教程后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考