AWS SQS未按预期触发Lambda与Step Functions,消息处理不全
问题描述
给S3存储桶配置SQS触发后,上传21个文件本该触发21次Step Functions执行,但实际仅触发14次,缺失7次事件。流程为:S3上传文件后推送消息至SQS,SQS触发Lambda执行;Lambda并发数设为5,理论上应按5个一批次完成21次Step Functions执行,但实际仅14次依次运行。Lambda代码如下:
import json import boto3 import os import uuid import time client = boto3.client('stepfunctions') ipparam ={} delay_seconds = 5 max_retry = 3 def lambda_handler(event, context): print('event',event) record = event['Records'][0] body = record['body'] s3EventBody= json.loads(body) s3EventBody= s3EventBody['Records'] bucket = s3EventBody[0]['s3']['bucket']['name'] key = s3EventBody[0]['s3']['object']['key'] ipparam["bucket"] = bucket ipparam["key"] = key print('bucket',bucket) print('key',key) transactionid = str(uuid.uuid4()) for retry_attempt in range(1,max_retry+1): try: response = client.start_execution( stateMachineArn=os.environ['stepFunctionArn'], name = transactionid, input=json.dumps(ipparam) ) print('responssee',response) executionArn = response['executionArn'] break except Exception as e: print("Error Invoking step function : ",e) if retry_attempt < max_retry: print('retrying in seconds : ',delay_seconds ) time.sleep(delay_seconds) else: print('Max retries exceeded. Please check logs for further details') raise for i in range(1,15): response = client.describe_execution( executionArn=executionArn ) execution_status = response['status'] print('execution_status',execution_status) if execution_status in ('SUCCEEDED','FAILED'): return { 'statusCode': 200, 'body': execution_status } time.sleep(60)
问题原因分析
- 未处理SQS批量消息:Lambda触发时,SQS可能批量发送多条消息到
event['Records']中,但当前代码仅处理第一条消息,处理完成后直接返回,SQS会将这批所有消息标记为已处理并删除,导致剩余消息丢失,最终部分S3上传事件未触发Step Functions。 - Lambda执行超时风险:代码最后通过循环等待Step Functions执行完成,单次循环最长等待14分钟,若Lambda配置的超时时间小于该时长,会直接触发超时。Lambda超时后,SQS会将消息放回队列重试,多次失败后消息可能进入死信队列,导致事件丢失。
- S3事件批量处理遗漏:S3推送至SQS的消息中,
s3EventBody['Records']可能包含多个文件上传事件,但代码仅处理第一条,同样会导致部分事件遗漏。
修复方案
- 遍历处理所有SQS消息:循环处理
event['Records']中的每一条SQS消息,避免批量消息丢失。 - 处理S3批量事件:针对S3可能批量推送的多文件事件,遍历
s3EventBody['Records']中的每一条记录。 - 移除Step Functions等待逻辑:Lambda仅负责启动Step Functions执行,无需等待其运行完成,大幅缩短Lambda执行时间,避免超时。
- 优化错误处理:单条消息处理失败时不中断整体流程,同时保留合理的重试机制,确保Step Functions启动成功率。
- 检查SQS配置:确认SQS的可见性超时大于Lambda超时时间,避免消息被重复投递;同时查看死信队列,确认是否有丢失的消息进入其中。
修改后的Lambda代码
import json import boto3 import os import uuid import time client = boto3.client('stepfunctions') delay_seconds = 5 max_retry = 3 def lambda_handler(event, context): # 遍历所有SQS消息 for record in event['Records']: try: body = record['body'] s3_event_body = json.loads(body) # 遍历S3事件中的每一条文件记录 for s3_record in s3_event_body['Records']: bucket = s3_record['s3']['bucket']['name'] key = s3_record['s3']['object']['key'] params = {"bucket": bucket, "key": key} transaction_id = str(uuid.uuid4()) # 重试启动Step Functions for retry_attempt in range(1, max_retry + 1): try: response = client.start_execution( stateMachineArn=os.environ['stepFunctionArn'], name=transaction_id, input=json.dumps(params) ) print(f"Step Functions启动成功,ARN: {response['executionArn']}") break except Exception as e: print(f"第{retry_attempt}次启动失败: {str(e)}") if retry_attempt < max_retry: time.sleep(delay_seconds) else: print(f"重试耗尽,无法启动Step Functions,Bucket: {bucket}, Key: {key}") raise except Exception as e: print(f"SQS消息处理失败: {str(e)}") continue return { 'statusCode': 200, 'body': "所有消息处理完成" }
内容的提问来源于stack exchange,提问作者VKRV
相关产品推荐
相关产品推荐

