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

如何用Python版AWS Lambda逐个处理S3桶文件,应对中途新增直至全处理完毕

基于AWS Lambda的S3文件串行处理与Glue ETL触发方案

需求要点:Appflow持续向S3指定路径推送文件,需串行触发Glue ETL(不可并发),确保所有文件(包括处理过程中新增的)都被处理,完成后归档/删除文件。
核心解决逻辑:用分布式锁避免多Lambda并发执行,处理完当前批次后重新扫描S3目录,循环直至无待处理文件。


1. 前置资源准备

  • 创建DynamoDB表(如FileProcessingLock):主键为LockKey(固定字符串值如S3ProcessingLock),包含IsLocked(布尔型)、LastUpdated(时间戳)字段,开启TTL自动释放锁(防止Lambda崩溃导致锁永久占用)。
  • 配置Lambda权限:允许访问S3(读、写、移动)、DynamoDB(读写)、Glue(触发工作流、查询运行状态)。
  • 配置S3事件触发:当有新文件上传时触发Lambda(事件批量不影响串行逻辑,锁机制会自动过滤并发请求)。

2. 核心Python代码实现

以下代码包含锁管理、文件扫描、Glue触发、文件归档全流程:

import boto3
import time
from botocore.exceptions import ClientError

# 初始化AWS客户端
s3 = boto3.client('s3')
dynamodb = boto3.resource('dynamodb')
glue = boto3.client('glue')

# 自定义配置参数
CONFIG = {
    "BUCKET_NAME": "your-target-bucket",
    "PROCESS_PREFIX": "appflow/incoming/",
    "ARCHIVE_PREFIX": "appflow/archived/",
    "LOCK_TABLE": "FileProcessingLock",
    "LOCK_KEY": "S3ProcessingLock",
    "GLUE_WORKFLOW_NAME": "your-glue-workflow",
    "MAX_RETRIES": 3,
    "POLL_INTERVAL": 30  # 轮询Glue工作流状态的间隔(秒)
}

def acquire_lock():
    """获取分布式锁,成功返回True,失败返回False"""
    table = dynamodb.Table(CONFIG["LOCK_TABLE"])
    try:
        table.put_item(
            Item={
                'LockKey': CONFIG["LOCK_KEY"],
                'IsLocked': True,
                'LastUpdated': int(time.time())
            },
            ConditionExpression='attribute_not_exists(LockKey) OR IsLocked = :false',
            ExpressionAttributeValues={':false': False},
            ReturnValues='ALL_OLD'
        )
        return True
    except ClientError as e:
        if e.response['Error']['Code'] == 'ConditionalCheckFailedException':
            # 锁已被其他进程占用
            return False
        raise

def release_lock():
    """释放分布式锁,确保异常场景下也能执行"""
    table = dynamodb.Table(CONFIG["LOCK_TABLE"])
    try:
        table.put_item(
            Item={
                'LockKey': CONFIG["LOCK_KEY"],
                'IsLocked': False,
                'LastUpdated': int(time.time())
            }
        )
    except Exception as e:
        print(f"释放锁失败: {str(e)}")

def get_pending_files():
    """获取S3待处理目录的文件列表,按上传时间升序排序"""
    try:
        response = s3.list_objects_v2(
            Bucket=CONFIG["BUCKET_NAME"],
            Prefix=CONFIG["PROCESS_PREFIX"],
            Delimiter='/'
        )
        if 'Contents' not in response:
            return []
        # 按文件上传时间排序,确保先处理早到的文件
        files = sorted(response['Contents'], key=lambda x: x['LastModified'])
        # 过滤掉目录本身
        return [f['Key'] for f in files if not f['Key'].endswith('/')]
    except ClientError as e:
        print(f"扫描S3文件失败: {str(e)}")
        return []

def trigger_glue_workflow(file_key):
    """触发Glue工作流并等待执行完成"""
    try:
        # 启动Glue工作流,传递文件路径作为参数
        start_response = glue.start_workflow_run(
            Name=CONFIG["GLUE_WORKFLOW_NAME"],
            RunProperties={
                's3_file_path': f"s3://{CONFIG['BUCKET_NAME']}/{file_key}"
            }
        )
        run_id = start_response['RunId']
        
        # 轮询工作流状态
        while True:
            run_status = glue.get_workflow_run(
                Name=CONFIG["GLUE_WORKFLOW_NAME"],
                RunId=run_id
            )['WorkflowRun']['Status']
            
            if run_status in ['SUCCEEDED', 'FAILED', 'STOPPED']:
                break
            time.sleep(CONFIG["POLL_INTERVAL"])
        
        if run_status != 'SUCCEEDED':
            raise Exception(f"Glue工作流执行失败,状态: {run_status}")
        return True
    except Exception as e:
        print(f"处理文件{file_key}失败: {str(e)}")
        return False

def archive_file(file_key):
    """将处理完成的文件移至归档目录并删除原文件"""
    try:
        # 构造归档路径
        archive_key = CONFIG["ARCHIVE_PREFIX"] + file_key.split(CONFIG["PROCESS_PREFIX"])[1]
        # 复制文件到归档目录
        s3.copy_object(
            Bucket=CONFIG["BUCKET_NAME"],
            CopySource=f"{CONFIG['BUCKET_NAME']}/{file_key}",
            Key=archive_key
        )
        # 删除原文件
        s3.delete_object(Bucket=CONFIG["BUCKET_NAME"], Key=file_key)
        print(f"文件{file_key}已归档至{archive_key}")
        return True
    except ClientError as e:
        print(f"归档文件{file_key}失败: {str(e)}")
        return False

def lambda_handler(event, context):
    # 尝试获取锁,获取失败则直接退出(已有处理进程在运行)
    if not acquire_lock():
        print("已有处理进程在运行,退出当前Lambda")
        return
    
    try:
        while True:
            pending_files = get_pending_files()
            if not pending_files:
                print("无待处理文件,结束处理")
                break
            
            # 逐个处理文件
            for file_key in pending_files:
                print(f"开始处理文件: {file_key}")
                if trigger_glue_workflow(file_key):
                    archive_file(file_key)
                else:
                    # 处理失败可根据需求移至错误目录或重试
                    print(f"文件{file_key}处理失败,跳过归档")
            
            # 处理完当前批次后重新扫描,覆盖处理期间新增的文件
            print("当前批次处理完成,重新扫描待处理文件")
    
    finally:
        # 无论成功失败,必须释放锁
        release_lock()
    
    return {"statusCode": 200, "body": "处理完成"}

3. 关键注意事项

  • Lambda超时设置:由于需要等待Glue工作流完成,需将Lambda超时时间设置为最大值(15分钟);若Glue工作流运行时间超过15分钟,建议改用AWS Step Functions编排,避免Lambda超时。
  • 锁的自动释放:在DynamoDB表中为LastUpdated字段开启TTL,设置锁过期时间(如30分钟),防止Lambda崩溃导致锁永久占用。
  • 错误处理:可根据实际需求扩展失败文件处理逻辑(如移至错误目录、触发告警)。
  • 权限配置:确保Lambda角色拥有以下权限:
    • S3:ListBucket、GetObject、PutObject、DeleteObject
    • DynamoDB:PutItem、GetItem
    • Glue:StartWorkflowRun、GetWorkflowRun

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 22:10:58