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

实现批量写入后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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 14:34:53