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

如何解决Python Lambda读取S3存储桶ZIP文件失败问题?

S3触发Lambda读取ZIP文件失败排查与修复

问题背景

开发后端功能实现App A向App B导入数据,流程如下:

  • 用户在网站上传包含CSV的ZIP压缩包
  • ZIP文件存入S3存储桶
  • 文件存入后触发Python Lambda函数
  • Lambda读取ZIP并处理其中的CSV数据

目前前3步正常,但第4步Lambda无法读取/处理文件。本地运行代码没问题,推测是S3对象获取方式错误导致read_zip失效。CloudWatch日志仅打印文件名,无后续处理日志。

原代码

import boto3
import csv
from io import TextIOWrapper, BytesIO
import zipfile
import requests
import re
from tqdm import tqdm
import concurrent.futures

s3Client = boto3.client('s3')

def lambda_handler(event,context):
    bucket = event['Records'][0]['s3']['bucket']['name']
    filename = event['Records'][0]['s3']['object']['key']
    # 打印过这两个值,显示正常

    # response = s3Client.get_object(Bucket=bucket, Key=filename)
    
    usefile = 'https://' + bucket + '.s3.ap-south-1.amazonaws.com/' + filename
    print(usefile)
    # 文件名打印正确

    def read_csv(file):
        to_return = []
        reader = csv.DictReader(TextIOWrapper(file, 'utf-8'))
        for row in reader:
            to_return.append(row)
        return to_return

    def read_zip(usefile):
        with zipfile.ZipFile(usefile, 'r') as APPA_file:
            with APPA_file.open("file1.csv", mode='r') as f:
                file1 = read_csv(f)
            with APPA_file.open("file2.csv", mode='r') as f:
                file2 = read_csv(f)
        return file1, file2

    def get_APPB_url(APPA_uri):
        resp = requests.get(APPA_uri)
        if resp.status_code != 200:
            return None
        # 提取APPB链接
        re_match = re.findall('href="(https://www.X.org/.+/)"', resp.text)
        if not re_match:
            return None
        print(resp.text)
        return re_match[0]

    def rate_on_APPB(APPB_url, rating):
        re_match = re.findall('X.org/(.+)/', APPB_url)
        if not re_match:
            return None
        
        APPB_id = re_match[0]
        req_body = {
            "query": "<query used>",
            "operationName": "<xyz>",
            "variables": "<variables>"
        }
        headers = {
            "content-type": "application/json",
            "x-hasura-admin-secret": "XXX"
        }

        resp = requests.post("XXX", json=req_body, headers=headers)

        if resp.status_code != 200:
            raise ValueError(f"Hasura query failed. Code: {resp.status_code}")
        else:
            print(APPB_id)

        json_resp = resp.json()
        if 'errors' in json_resp and len(json_resp['errors']) > 0:
            first_error_msg = json_resp['errors'][0]['message']
            if 'Authentication' in first_error_msg:
                print("Failed to authenticate with cookie")
                exit(1)
            else:
                raise ValueError(first_error_msg)

    def APPA_to_APPB(APPA_dict):
        APPB_url = get_APPB_url(APPA_dict['APPA URI'])
        if APPB_url is None:
            raise ValueError("Cannot find APPB title")
        rate_on_APPB(APPB_url, int(float(APPA_dict['Rating']) * 2))

    def main():
        file1, file2 = read_zip(usefile)
        success = []
        errors = []

        with tqdm(total=len(file1)) as pbar:
            with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
                future_to_url = {
                    executor.submit(APPA_to_APPB, APPA_dict): APPA_dict for APPA_dict in file1
                }
                try:
                    for future in concurrent.futures.as_completed(future_to_url):
                        APPA_dict = future_to_url[future]
                        pbar.update(1)
                        try:
                            success.append(future.result())
                        except Exception as e:
                            errors.append({"APPA_dict": APPA_dict, "error": e})
                except KeyboardInterrupt:
                    executor._threads.clear()
                    concurrent.futures.thread._threads_queues.clear()

        print(f"Successfully rated: {len(success)} ")
        print(f"{len(errors)} Errors")
        for error in errors:
            print(f"	{error['APPA_dict']['Name']} ({error['APPA_dict']['Year']}): {error['error']}")

    if __name__ == '__main__':
        main()

核心问题分析

  1. Lambda未执行main函数:Lambda运行时是调用lambda_handler函数,而非以脚本方式启动,所以if __name__ == '__main__': main()这段代码永远不会触发,这是日志无后续输出的根本原因。
  2. S3文件获取方式错误:直接把S3的HTTP URL传给ZipFile是无效的,ZipFile需要本地文件路径或文件流对象,不能直接处理远程URL。
  3. 缺乏错误捕获与日志:原代码没有对read_zip等关键步骤做异常捕获,即使出错也不会在CloudWatch中输出错误信息。

修复方案

1. 触发main函数执行

删除if __name__ == '__main__': main(),直接在lambda_handler末尾调用main(),同时将bucket和filename作为参数传入main。

2. 正确从S3读取ZIP文件

使用boto3.get_object获取S3对象的字节流,用BytesIO包装后传给ZipFile,修改read_zip函数:

def read_zip(bucket, filename):
    # 从S3获取对象字节流
    response = s3Client.get_object(Bucket=bucket, Key=filename)
    zip_bytes = response['Body'].read()
    # 用BytesIO包装成文件流
    with zipfile.ZipFile(BytesIO(zip_bytes), 'r') as APPA_file:
        with APPA_file.open("file1.csv", mode='r') as f:
            file1 = read_csv(f)
        with APPA_file.open("file2.csv", mode='r') as f:
            file2 = read_csv(f)
    return file1, file2

3. 添加异常捕获与日志

在核心函数中添加try-except块,捕获并打印异常,便于在CloudWatch中排查问题。

4. 优化代码结构

将嵌套函数移到lambda_handler外部,提升代码可读性和可维护性。

修复后完整代码示例

import boto3
import csv
from io import TextIOWrapper, BytesIO
import zipfile
import requests
import re
from tqdm import tqdm
import concurrent.futures

s3Client = boto3.client('s3')

def read_csv(file):
    to_return = []
    reader = csv.DictReader(TextIOWrapper(file, 'utf-8'))
    for row in reader:
        to_return.append(row)
    return to_return

def read_zip(bucket, filename):
    try:
        response = s3Client.get_object(Bucket=bucket, Key=filename)
        zip_bytes = response['Body'].read()
        with zipfile.ZipFile(BytesIO(zip_bytes), 'r') as APPA_file:
            with APPA_file.open("file1.csv", mode='r') as f:
                file1 = read_csv(f)
            with APPA_file.open("file2.csv", mode='r') as f:
                file2 = read_csv(f)
        print(f"Successfully read {len(file1)} rows from file1.csv")
        return file1, file2
    except Exception as e:
        print(f"Failed to read zip file: {str(e)}")
        raise e

def get_APPB_url(APPA_uri):
    try:
        resp = requests.get(APPA_uri)
        resp.raise_for_status()  # 抛出HTTP错误
        re_match = re.findall('href="(https://www.X.org/.+/)"', resp.text)
        if not re_match:
            print(f"No APPB URL found in {APPA_uri}")
            return None
        return re_match[0]
    except Exception as e:
        print(f"Error fetching APPB URL: {str(e)}")
        return None

def rate_on_APPB(APPB_url, rating):
    re_match = re.findall('X.org/(.+)/', APPB_url)
    if not re_match:
        print(f"Invalid APPB URL: {APPB_url}")
        return None
    
    APPB_id = re_match[0]
    req_body = {
        "query": "<query used>",
        "operationName": "<xyz>",
        "variables": "<variables>"
    }
    headers = {
        "content-type": "application/json",
        "x-hasura-admin-secret": "XXX"
    }

    try:
        resp = requests.post("XXX", json=req_body, headers=headers)
        resp.raise_for_status()
        print(f"Successfully rated APPB ID: {APPB_id}")
        
        json_resp = resp.json()
        if 'errors' in json_resp and len(json_resp['errors']) > 0:
            first_error_msg = json_resp['errors'][0]['message']
            if 'Authentication' in first_error_msg:
                print("Failed to authenticate with cookie")
                raise ValueError("Authentication failure")
            else:
                raise ValueError(first_error_msg)
    except Exception as e:
        print(f"Error rating APPB ID {APPB_id}: {str(e)}")
        raise e

def APPA_to_APPB(APPA_dict):
    APPB_url = get_APPB_url(APPA_dict['APPA URI'])
    if APPB_url is None:
        raise ValueError(f"Cannot find APPB title for {APPA_dict['Name']}")
    rate_on_APPB(APPB_url, int(float(APPA_dict['Rating']) * 2))

def main(bucket, filename):
    file1, file2 = read_zip(bucket, filename)
    success = []
    errors = []

    with tqdm(total=len(file1)) as pbar:
        with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
            future_to_url = {
                executor.submit(APPA_to_APPB, APPA_dict): APPA_dict for APPA_dict in file1
            }
            for future in concurrent.futures.as_completed(future_to_url):
                APPA_dict = future_to_url[future]
                pbar.update(1)
                try:
                    future.result()
                    success.append(APPA_dict)
                except Exception as e:
                    errors.append({"APPA_dict": APPA_dict, "error": str(e)})

    print(f"Successfully rated: {len(success)} entries")
    print(f"Failed entries: {len(errors)}")
    for error in errors:
        print(f"	{error['APPA_dict']['Name']} ({error['APPA_dict']['Year']}): {error['error']}")

def lambda_handler(event,context):
    try:
        bucket = event['Records'][0]['s3']['bucket']['name']
        filename = event['Records'][0]['s3']['object']['key']
        print(f"Starting processing for file: {filename} in bucket: {bucket}")
        
        main(bucket, filename)
        
        print("Processing completed successfully")
    except Exception as e:
        print(f"Processing failed with error: {str(e)}")
        raise e

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 04:45:35