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

AWS Lambda触发SQS时无法执行Step Function的问题排查

问题:Lambda添加SQS触发器后无法启动Step Function

手动运行Lambda时,从SQS读取消息、解码、启动Step Function的全流程均正常执行,但为Lambda添加SQS触发器后,虽然Lambda监控显示函数已被正确调用,却完全无法启动Step Function。已单独拆分启动Step Function的逻辑(手动运行正常),且为IAM角色授予了AdministratorAccess权限以确保权限充足。

Lambda代码如下:

import logging
import boto3
import base64
import json
import os
import ast

logger = logging.getLogger()
logger.setLevel(logging.INFO)


def get_queue_url(queue_name):
    """
    Get the URL of the queue dynamically based on its name.
    """
    sqs = boto3.client('sqs')
    response = sqs.get_queue_url(QueueName=queue_name)
    return response['QueueUrl']


def decode_celery_message(body):
    """
    Decode and parse the Celery message.
    """
    decoded_body = base64.b64decode(body).decode('utf-8')
    try:
        parsed_body = json.loads(decoded_body)
        if 'kwargsrepr' in parsed_body['headers']:
            kwargsrepr = parsed_body['headers']['kwargsrepr']
            task_id = parsed_body['headers']['id']
            kwargsrepr_dict = ast.literal_eval(kwargsrepr)
            kwargsrepr_dict['task_id'] = task_id
            return kwargsrepr_dict
        else:
            logger.error("'kwargsrepr' not found in the message headers.")
            return None
    except json.JSONDecodeError as e:
        logger.error("Error decoding JSON: %s", e)
        return None


def start_step_function_execution(parsed_message):
    """
    Start the execution of the Step Function.
    """
    stepfunctions_client = boto3.client('stepfunctions')
    logger.info(f"Starting Step Function execution with input: {parsed_message}")
    stepfunctions_client.start_execution(
        stateMachineArn=os.environ['STEP_FUNCTION_ARN'],
        input=json.dumps(parsed_message) 
    )
    logger.info("Step Function execution started successfully.")
    return parsed_message


def display_message_available(event, context):
    queue_name = os.environ["QUEUE_NAME"]  # Queue name
    queue_url = get_queue_url(queue_name)
    sqs = boto3.client('sqs')

    logger.info(f"Event details: {event}")
    
    all_messages = []  # List to store all decoded messages with commands
    
    while True:
        # Receive messages from the queue for processing
        response = sqs.receive_message(
            QueueUrl=queue_url,
        )
        
        # Check if any messages were received
        if 'Messages' in response:
            for message in response['Messages']:
                body = message['Body']
                receipt_handle = message['ReceiptHandle']
                
                # Decode and parse the message, get the command
                parsed_message = decode_celery_message(body)
                if parsed_message:
                    all_messages.append((parsed_message))
                    
                
                # Delete the message from the queue after processing
                sqs.delete_message(
                    QueueUrl=queue_url,
                    ReceiptHandle=receipt_handle
                )
                logger.info(f'Message {body} deleted from the {queue_name} after processing.')
                # Start Step Function execution
                start_step_function_execution(parsed_message)
                logger.info(f"Step Function execution started with input: {parsed_message}")
        else:
            logger.info("No more messages available in the queue.")
            break  # Exit the loop if no more messages are available
        

    return all_messages


# This is necessary for Lambda to know what function to invoke
lambda_handler = display_message_available

问题原因

  1. 触发器消息被忽略+重复处理冲突:SQS触发器会自动将消息推送到Lambda的event参数中,但你的代码完全没处理event里的消息,反而主动调用sqs.receive_message()拉取队列消息,同时还调用sqs.delete_message()删除消息——这会导致触发器推送的目标消息被提前删除或完全被忽略,后续启动Step Function的逻辑拿到的可能是无效数据,甚至根本没有可处理的消息。
  2. 异常无捕获,问题不可见:启动Step Function的逻辑没有异常捕获机制,当parsed_message为None(解码失败)或调用start_execution出错时,异常会被静默吞掉,没有详细日志输出,无法定位问题。
  3. 消息处理逻辑冗余:SQS触发器本身会处理消息的可见性和自动删除(Lambda执行成功则自动删消息,失败则重新推送),代码中手动调用delete_message属于重复操作,容易导致消息丢失或重复处理。

解决方法

1. 修改核心逻辑,处理触发器传入的消息

移除主动拉取SQS消息的代码,直接从event中读取触发器推送的消息,简化后的核心函数如下:

def display_message_available(event, context):
    logger.info(f"Event details: {event}")
    
    all_messages = []
    
    # 遍历触发器传入的消息记录
    for record in event['Records']:
        body = record['body']
        
        # 解码解析消息
        parsed_message = decode_celery_message(body)
        if parsed_message:
            all_messages.append(parsed_message)
            # 启动Step Function并捕获异常
            try:
                start_step_function_execution(parsed_message)
                logger.info(f"Step Function execution started with input: {parsed_message}")
            except Exception as e:
                logger.error(f"Failed to start Step Function: {str(e)}", exc_info=True)
                # 抛出异常让Lambda重试,避免无效消息被自动删除
                raise e
        else:
            logger.error(f"Failed to parse message body: {body}")
            # 解析失败可根据业务选择重试或丢弃,这里选择抛出异常重试
            raise ValueError("Invalid message format, missing required fields")
    
    return all_messages

2. 完善Step Function启动逻辑的异常捕获

在start_step_function_execution中添加异常处理,记录详细错误信息:

def start_step_function_execution(parsed_message):
    stepfunctions_client = boto3.client('stepfunctions')
    logger.info(f"Starting Step Function execution with input: {parsed_message}")
    try:
        stepfunctions_client.start_execution(
            stateMachineArn=os.environ['STEP_FUNCTION_ARN'],
            input=json.dumps(parsed_message) 
        )
        logger.info("Step Function execution started successfully.")
        return parsed_message
    except Exception as e:
        logger.error(f"Error starting Step Function execution: {str(e)}", exc_info=True)
        raise

3. 验证触发器配置

  • 确认SQS触发器的批量大小设置合理,避免一次推送过多消息导致Lambda超时
  • 确认触发器的可见性超时大于Lambda的超时时间,防止消息在Lambda处理完成前被重新推送
  • 确认SQS队列与Lambda处于同一AWS区域,避免跨区域调用问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 21:40:15