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

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']可能包含多个文件上传事件,但代码仅处理第一条,同样会导致部分事件遗漏。
修复方案
  1. 遍历处理所有SQS消息:循环处理event['Records']中的每一条SQS消息,避免批量消息丢失。
  2. 处理S3批量事件:针对S3可能批量推送的多文件事件,遍历s3EventBody['Records']中的每一条记录。
  3. 移除Step Functions等待逻辑:Lambda仅负责启动Step Functions执行,无需等待其运行完成,大幅缩短Lambda执行时间,避免超时。
  4. 优化错误处理:单条消息处理失败时不中断整体流程,同时保留合理的重试机制,确保Step Functions启动成功率。
  5. 检查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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 14:54:52