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桶名不区分大小写,但配置需严格匹配)。
- 确认S3桶的事件通知是否正确关联到该Lambda(因配置
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表的GSI
emailIndex状态是否为ACTIVE,确保查询逻辑能正常执行。
6. 其他排查点
- 检查Lambda执行角色的权限,通过IAM模拟器验证
sqs:SendMessage、dynamodb:PutItem等权限是否正常。 - 查看Lambda函数的并发限制,确认是否因并发不足导致函数无法执行。
- 核对所有环境变量值是否正确注入到Lambda函数中(在AWS控制台查看函数配置)。
内容的提问来源于stack exchange,提问作者Darla David
相关产品推荐
相关产品推荐

