如何用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
- S3:
内容的提问来源于stack exchange,提问作者Healer77Om
相关产品推荐
相关产品推荐

