如何基于AWS Step Functions实现SQS父子消息等待聚合逻辑?
基于AWS Step Functions实现无数据库的SQS消息聚合逻辑
一、Step Functions适配性说明
Step Functions完全适合该架构,它天生擅长异步状态跟踪、流程编排和等待型任务处理,刚好匹配你需要聚合关联消息、超时控制的核心需求,无需额外数据库即可实现状态持久化和事件驱动的流程推进。
二、无数据库的核心实现方案
整体思路
通过Step Functions执行名称绑定typeA消息ID,利用waitForTaskToken回调模式实现异步等待,配合SQS临时存储回调令牌,完成typeA与typeB消息的关联聚合:
1. 消息路由前置Lambda
所有SQS消息先经过一个路由Lambda,判断消息类型后触发对应逻辑:
typeA消息:调用Step Functions的StartExecutionAPI,执行名称设为typeA-<typeA-ID>,输入参数包含typeAId、expectedTypeBCount,同时设置执行超时为3600秒(1小时)。typeB消息:提取parentTypeAId,查询对应typeA执行的回调令牌,触发Step Functions的SendTaskSuccessAPI传递typeB数据。
2. TypeA触发的Step Functions状态机
状态机负责初始化聚合状态、等待typeB回调、检查聚合完成度,核心定义如下:
{ "Comment": "TypeA消息触发的TypeB聚合流程", "StartAt": "InitializeAggregation", "States": { "InitializeAggregation": { "Type": "Pass", "Parameters": { "typeAId.$": "$.typeAId", "expectedCount.$": "$.expectedTypeBCount", "receivedTypeBs": [] }, "Next": "StoreTaskToken" }, "StoreTaskToken": { "Type": "Task", "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:store-task-token", "Parameters": { "typeAId.$": "$.typeAId", "taskToken.$": "$$.Task.Token" }, "Next": "WaitForTypeBCallbacks" }, "WaitForTypeBCallbacks": { "Type": "Task", "Resource": "arn:aws:states:::lambda:invoke.waitForTaskToken", "Parameters": { "FunctionName": "dummy-placeholder-function", "Payload": { "taskToken.$": "$$.Task.Token" } }, "TimeoutSeconds": 3600, "Next": "MergeTypeBData", "Catch": [ { "ErrorEquals": ["States.Timeout"], "Result": "TIMEOUT: 1小时内未集齐所有TypeB消息", "End": true } ] }, "MergeTypeBData": { "Type": "Task", "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:deduplicate-and-merge-typeb", "Parameters": { "currentState.$": "$", "newTypeB.$": "$.Payload" }, "Next": "CheckAggregationStatus" }, "CheckAggregationStatus": { "Type": "Choice", "Choices": [ { "Variable": "$.receivedTypeBs.length", "NumericEqualsPath": "$.expectedCount", "Next": "RunCalculationLambda" } ], "Default": "WaitForTypeBCallbacks" }, "RunCalculationLambda": { "Type": "Task", "Resource": "arn:aws:lambda:REGION:ACCOUNT_ID:function:final-calculation-lambda", "Parameters": { "typeAId.$": "$.typeAId", "typeBData.$": "$.receivedTypeBs" }, "End": true } } }
3. 关键组件细节
- 令牌存储Lambda:将
typeAId和对应的TaskToken写入一个专用SQS队列,队列消息过期时间设为1小时,自动清理超时无效的令牌。 - 去重合并Lambda:每次收到
typeB消息时,检查receivedTypeBs数组中是否已存在当前typeB的ID,避免SQS至少一次投递导致的重复计数,合并后更新状态机的聚合数据。 - TypeB消息处理Lambda:提取
parentTypeAId后,从专用SQS队列中查询对应TaskToken,调用SendTaskSuccess将typeB数据传递给状态机,触发聚合检查。
三、核心逻辑说明
typeA消息触发后:状态机初始化聚合状态,存储回调令牌,进入等待回调状态,直到超时或收到足够typeB消息。typeB消息触发后:通过令牌找到对应typeA的状态机执行,传递typeB数据,状态机合并数据并检查是否集齐,未集齐则继续等待,集齐则调用计算Lambda。- 超时处理:状态机内置1小时超时,超时后自动返回指定状态,无需额外监控。
内容的提问来源于stack exchange,提问作者Maria Adamsky
相关产品推荐
相关产品推荐

