实现批量写入后Lambda函数未执行且无报错排查求助
问题
原本配置15分钟超时、128MB内存的Lambda函数,可正常将.xlsx文件数据行插入DynamoDB。现在要处理含16万行数据的特定文件,已将内存调整至1024MB并实现批量写入请求,触发器、权限等其他配置未变更且确认正常,但函数完全未触发,CloudWatch中无任何日志生成。
函数代码
import os import boto3 import pandas as pd from io import BytesIO import logging from datetime import datetime from concurrent.futures import ThreadPoolExecutor # 配置日志 logger = logging.getLogger() logger.setLevel(logging.INFO) # AWS服务客户端 s3 = boto3.client('s3') dynamodb = boto3.client('dynamodb') # 定义S3桶名和DynamoDB表名 # BUCKET_NAME = 'sellersinternal' TABLE_NAME = 'Seller' # 定义文件名到列映射的对应关系 # 可根据其他文件添加更多映射(当前基于电商平台爬取数据) file_column_mappings = { 'sellers_testing.xlsx': { 'Seller Name': ['SellerName'], 'Listing URL': ['SellerUrl'], 'Headquarter': ['SellerHq'], 'Revenue': ['SellerRevenue'], 'Year Founded': ['SellerYear'], 'Number of Employees': ['SellerEmployeeCount'], 'Icon URL': ['SellerOriginalIconUrl'], # 新增列 'Category L2': ['SellerCategory'], 'Category L3': ['SellerSubCategory'] } } def batch_write_items(table_name, items): dynamodb.batch_write_item(RequestItems={table_name: items}) def lambda_handler(event, context): logger.info("Lambda函数开始执行...") # 从S3事件中获取上传的文件 file_obj = event['Records'][0] bucket_name = file_obj['s3']['bucket']['name'] file_key = file_obj['s3']['object']['key'] # 根据文件名获取对应的列映射 column_mapping = file_column_mappings.get(os.path.basename(file_key), None) try: # 从S3读取XLSX文件 response = s3.get_object(Bucket=bucket_name, Key=file_key) excel_data = response['Body'].read() # 用Pandas解析XLSX文件 excel_df = pd.read_excel(BytesIO(excel_data)) # 获取S3文件上传时间 s3_upload_time = file_obj['eventTime'] # 初始化批量处理的项目列表 batch_items = [] # 遍历每一行,准备插入DynamoDB的项目 for _, row in excel_df.iterrows(): item = { 'SellerTimestamp': {'S': str(datetime.now())} } for excel_column, dynamodb_attributes in column_mapping.items(): for dynamodb_attribute in dynamodb_attributes: if dynamodb_attribute not in item: item[dynamodb_attribute] = {'S': str(row[excel_column])} else: if not isinstance(item[dynamodb_attribute], list): item[dynamodb_attribute] = [item[dynamodb_attribute]] item[dynamodb_attribute].append({'S': str(row[excel_column])}) if 'SellerName' in item: item['SellerNameLC'] = {'S': item['SellerName']['S'].lower()} seller_name = item.get('SellerName', {}).get('S', '') seller_name_lc = ''.join(seller_name.split()).lower() # 移除空格并转为小写 item['SellerId'] = {'S': seller_name_lc} # 将项目添加到批量列表 batch_items.append({'PutRequest': {'Item': item}}) # 当批量列表达到25条时,并行发起批量写入请求 if len(batch_items) == 25: with ThreadPoolExecutor(max_workers=15) as executor: futures = [executor.submit(batch_write_items, TABLE_NAME, batch_items)] batch_items = [] # 处理剩余的批量项目 if batch_items: batch_write_items(TABLE_NAME, batch_items) except Exception as e: logger.error("错误: %s", e) return { 'statusCode': 500, 'body': '错误: ' + str(e) }
排查与解决方案
一、触发器与调用链路排查
- 确认S3触发器配置:检查事件类型是否为
s3:ObjectCreated:*,前缀/后缀规则是否匹配目标文件名,触发器状态是否为启用,目标Lambda ARN是否正确 - 查看CloudTrail日志:搜索
lambda:InvokeFunction事件,确认S3是否发起了Lambda调用请求,若存在调用被拒绝记录,需检查权限配置 - 检查Lambda并发配额:在Lambda控制台“配置>并发”页面,确认账户或函数的并发数未耗尽,避免新请求被丢弃
二、函数初始化与加载问题
- 验证依赖兼容性:1024MB内存环境下,重新打包Pandas等依赖包,确保与Lambda运行时环境兼容;若使用层,确认层未损坏
- 处理大文件内存溢出:16万行Excel文件体积可能过大,改用分块读取避免内存不足:
# 替换原读取代码为分块读取 excel_df = pd.read_excel(BytesIO(excel_data), chunksize=1000) for chunk in excel_df: for _, row in chunk.iterrows(): # 保留原行处理逻辑 - 增加列映射校验:在
column_mapping赋值后添加判断,避免因文件名不匹配导致后续逻辑崩溃:if not column_mapping: logger.error(f"未找到文件{os.path.basename(file_key)}的列映射") return {'statusCode': 400, 'body': '无效的文件名'}
三、日志与权限校验
- 确认Lambda执行角色权限:检查角色是否拥有
logs:CreateLogGroup、logs:CreateLogStream、logs:PutLogEvents权限,避免日志无法生成 - 检查CloudWatch日志组:确认对应Lambda的日志组存在,若日志流未创建,说明函数未被触发或初始化阶段直接崩溃
四、批量写入逻辑优化
- 优化线程池使用:避免每次批量都创建新的ThreadPoolExecutor,将线程池初始化移至handler外部,复用线程资源
- 增加批量写入错误处理:捕获
batch_write_item返回的未处理项目,进行重试,避免数据丢失
内容的提问来源于stack exchange,提问作者Jam Filho
相关产品推荐
相关产品推荐

