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

Serverless部署AWS Lambda异常:SQS队列空、DynamoDB未更新

问题描述

我用Serverless Framework部署了AWS Lambda函数,流程是S3上传CSV文件后触发Lambda,将CSV每行数据发送到SQS队列,再由另一个Lambda监听SQS,将数据写入DynamoDB表。部署成功,日志未显示异常,但SQS队列始终为空,DynamoDB也没有数据写入。目前完成了前4步,卡在第5步(SQS到DynamoDB的环节,但实际SQS都没数据)。

serverless.yml

service: challenge1
frameworkVersion: '3'

provider:
  name: aws
  runtime: python3.8
  lambdaHashingVersion: '20201221'
  iamRoleStatements:
    - Effect: "Allow"
      Action: "dynamodb:*"
      Resource: "*"
    - Effect: "Allow"
      Action: "apigateway:*"
      Resource: "*"
    - Effect: "Allow"
      Action: "s3:*"
      Resource: "*"
    - Effect: "Allow"
      Action: "sqs:*"
      Resource: "*"
  environment:
    DYNAMODB_CARDS_TABLE_NAME: challenge1 
    S3_BUCKETNAME: serverlesschallenge-darla
    QUEUE_URL: https://sqs.us-east-1.amazonaws.com/874957933250/serverlesschallenge-darla

functions:
  prepareSQSjobS3:
    handler: handler.prepare_sqs_job
    events:
      - s3:
          bucket: serverlesschallenge-darla
          event: s3:ObjectCreated:Put
          existing: true
          rules:
            - suffix: .csv
  prepareSQSjobSQS:
    handler: handler.process_sqs_job
    events:
      - sqs:
          arn: "arn:aws:sqs:us-east-1:874957933250:serverlesschallenge-darla"

package:
  exclude:
    - venv/**
    - node_modules/**

resources:
  Resources:
    LoyaltyCardDynamodbTable:
      Type: 'AWS::DynamoDB::Table'
      Properties:
        AttributeDefinitions:
          - AttributeName: card_number
            AttributeType: S
          - AttributeName: email
            AttributeType: S
        KeySchema:
          - AttributeName: card_number
            KeyType: HASH
        BillingMode: PAY_PER_REQUEST
        TableName: ${self:provider.environment.DYNAMODB_CARDS_TABLE_NAME}
        GlobalSecondaryIndexes:
            - IndexName: emailIndex
              KeySchema:
                - AttributeName: email
                  KeyType: HASH
              Projection:
                ProjectionType: ALL

plugins:
  - serverless-python-requirements

handler.py

import json
import string
import random
import os
import boto3
import urllib.parse
import csv
import sys
from io import StringIO

from dynamodb_gateway import DynamodbGateway

s3 = boto3.client('s3')
sqs = boto3.client('sqs')
queue_url = os.getenv('QUEUE_URL')

#aws lambda trigger when theres new s3 file. reads line by line
def prepare_sqs_job(event, context):
    try:
        print(f"Received S3 event: {json.dumps(event)}")

        bucket_name = os.getenv("S3_BUCKETNAME")

        # Get the object details from the S3 event
        s3_record = event['Records'][0]['s3']
        bucket = s3_record['bucket']['name']
        file_key = urllib.parse.unquote_plus(s3_record['object']['key'], encoding='utf-8')

        # Download the file from S3
        response = s3.get_object(Bucket=bucket, Key=file_key)
        file_content = response['Body'].read().decode('utf-8')
        print(f"Object uploaded: s3://{bucket}/{file_key}")

        # Process CSV file and send each row as a message to SQS
        rows = [row for i, row in enumerate(csv.reader(StringIO(file_content))) if i > 0]

        message_attrs = {'AttributeName': {'StringValue': 'AttributeValue', 'DataType': 'String'}}
        for row in rows:
            print(row)
            sqs.send_message(
                QueueUrl=queue_url,
                MessageBody=row[0],
                MessageAttributes=message_attrs,
            )

        message = 'Messages accepted!'
        print(message)
        response = {"statusCode": 200, "body": json.dumps({"status": "success", "message": message})}

    except Exception as e:
        print(f'Error: {str(e)}')
        response = {"statusCode": 500, "body": json.dumps({"status": "error", "message": str(e)})}

    return response

def process_sqs_job(event, context):
    try:
        print(f"Received SQS event: {json.dumps(event)}")

        table_name = os.getenv("DYNAMODB_CARDS_TABLE_NAME")

        for record in event['Records']:
            # Parse JSON content from SQS message
            message_body = json.loads(record['body'])

            if isinstance(message_body, dict):
                # Extract necessary information from the message
                card_number = message_body.get('card_number')
                first_name = message_body.get('first_name')
                last_name = message_body.get('last_name')
                email = message_body.get('email')
                points = message_body.get('points')

                # Check if the email already exists in the DynamoDB table
                if email_exists(table_name, email):
                    print(f"Email {email} already used. Skipping...")
                    continue

                # Create a loyalty card in DynamoDB
                loyalty_card = {
                    "card_number": card_number,
                    "first_name": first_name,
                    "last_name": last_name,
                    "email": email,
                    "points": points
                }

                DynamodbGateway.upsert(
                    table_name=table_name,
                    mapping_data=[loyalty_card],
                    primary_keys=["card_number"]
                )

                print(f"Loyalty card created: {loyalty_card}")

        message = 'Messages processed successfully!'
        print(message)
        response = {"statusCode": 200, "body": json.dumps({"status": "success", "message": message})}

    except Exception as e:
        print(f'Error: {str(e)}')
        response = {"statusCode": 500, "body": json.dumps({"status": "error", "message": str(e)})}

    return response


def email_exists(table_name, email):
    # Check if the email already exists in the DynamoDB table using GSI
    result = DynamodbGateway.query_index_by_partition_key(
        index_name="emailIndex",
        table_name=table_name,
        partition_key_name="email",
        partition_key_query_value=email
    )

    return bool(result)

已尝试添加send_message的异常捕获并打印结果,同时查看了两个函数的CloudWatch日志,但未发现异常。


排查思路与解决建议

1. 确认S3触发的Lambda是否实际执行

  • 检查prepare_sqs_job的CloudWatch日志,确认是否输出Received S3 event和Object uploaded,验证S3触发逻辑是否生效。
  • 查看日志中打印的row内容,确认CSV解析是否正确,每行数据是否符合预期格式。
  • 如果S3触发未执行:
    • 确认S3桶的事件通知是否正确关联到该Lambda(因配置existing: true,需手动检查桶的事件通知配置)。
    • 核对桶名与serverless.yml中配置的bucket是否完全一致(S3桶名不区分大小写,但配置需严格匹配)。

2. 修复SQS消息发送的格式错误

目前prepare_sqs_job中发送的MessageBody=row[0]仅传递了CSV每行的第一列,而process_sqs_job尝试用json.loads将其解析为字典,这会直接抛出异常(但被外层try捕获未详细打印)。需修改消息发送逻辑:

# 在prepare_sqs_job的循环中,将整行数据构造为字典后转成JSON发送
for row in rows:
    print(row)
    # 假设CSV列顺序为card_number,first_name,last_name,email,points
    message_body = {
        "card_number": row[0],
        "first_name": row[1],
        "last_name": row[2],
        "email": row[3],
        "points": row[4]
    }
    try:
        resp = sqs.send_message(
            QueueUrl=queue_url,
            MessageBody=json.dumps(message_body),
            MessageAttributes=message_attrs,
        )
        print(f"Sent message ID: {resp['MessageId']}")
    except Exception as send_err:
        print(f"Failed to send message: {str(send_err)}")

3. 验证SQS队列与触发器配置

  • 核对QUEUE_URL环境变量与AWS控制台中SQS队列的实际URL是否一致。
  • 检查SQS队列的触发器配置,确认已关联到prepareSQSjobSQS函数。
  • 查看SQS监控指标NumberOfMessagesSent,确认是否有消息被成功发送到队列。

4. 修复SQS消息处理的逻辑漏洞

在process_sqs_job中,当消息体不是字典时,代码会直接跳过处理且无日志输出,需添加日志排查:

message_body = json.loads(record['body'])
if isinstance(message_body, dict):
    # 原有处理逻辑
else:
    print(f"Invalid message format, not a dict: {message_body}")

5. 验证DynamoDB写入逻辑

  • 替换自定义DynamodbGateway为原生boto3代码测试,确认写入逻辑是否正常:
# 替换DynamodbGateway.upsert的调用
dynamodb = boto3.resource('dynamodb')
table = dynamodb.Table(table_name)
table.put_item(Item=loyalty_card)
  • 检查DynamoDB表的GSIemailIndex状态是否为ACTIVE,确保查询逻辑能正常执行。

6. 其他排查点

  • 检查Lambda执行角色的权限,通过IAM模拟器验证sqs:SendMessage、dynamodb:PutItem等权限是否正常。
  • 查看Lambda函数的并发限制,确认是否因并发不足导致函数无法执行。
  • 核对所有环境变量值是否正确注入到Lambda函数中(在AWS控制台查看函数配置)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 04:37:05