如何优化AWS Lambda函数实现S3多文件上传后仅单次运行Glue作业
S3触发Glue作业重复运行优化方案
问题背景
现有Lambda函数监听S3新对象创建事件,每收到1个新文件事件就启动1次Glue作业,Kafka Connect按批次写入S3时(单批5000条记录,或空闲超时强制刷写),单批次会产生多个新文件,导致Glue作业重复触发浪费资源。
原代码还存在1个显性bug:print(bucket, fileName)语句中bucket变量未定义,实际运行会抛错。
import json import boto3 from urllib.parse import unquote_plus def lambda_handler(event, context): bucketName = event["Records"][0]["s3"]["bucket"]["name"] fileNameFull = event["Records"][0]["s3"]["object"]["key"] fileName = unquote_plus(fileNameFull) print(bucket, fileName) glue = boto3.client('glue') response = glue.start_job_run( JobName = 'My_Job_Glue', Arguments = { '--s3_target_path_key': fileName, '--s3_target_path_bucket': bucketName } ) return { 'statusCode': 200, 'body': json.dumps('Hello from Lambda!') }
核心优化思路
不要让S3事件直接触发作业启动,通过事件缓冲+批量聚合+运行状态互斥的逻辑,等单批次文件全部写入完成后,只启动1次Glue作业处理全量文件。
具体实现步骤
- 调整事件链路:S3新对象创建事件不要直接投递Lambda,先投递到SQS队列做缓冲。队列参数配置:
- 消息可见性超时 = Kafka Connect最大空闲刷写时间 + 2分钟(比如Kafka空闲刷写阈值是5分钟,就设为7分钟)
- 消息保留期设为1天即可
- 配置Lambda触发源为上述SQS队列:
- 批处理大小设为20(覆盖单批次Kafka写入的最大文件拆分数量即可)
- 批处理窗口设为60秒,攒齐窗口内的所有事件再触发Lambda
- 改造Lambda逻辑:
- 遍历拉取到的所有事件,过滤Kafka Connect写入过程中产生的
.tmp临时文件(文件写完才会重命名为正式后缀,临时文件不需要处理) - 对有效文件按存储桶做去重聚合,避免S3/SQS事件重复投递导致的重复统计
- 调用Glue接口检查目标作业是否存在
STARTING/RUNNING/WAITING/STOPPING状态的运行实例,有运行中实例时不启动新作业,仅把新文件加入待处理清单 - 无运行中实例时,把聚合后的待处理文件清单存到S3,将清单路径作为参数传给Glue启动作业
优化后代码示例
import json import boto3 from urllib.parse import unquote_plus from datetime import datetime # 初始化客户端放handler外,复用执行环境连接 glue = boto3.client('glue') s3 = boto3.client('s3') # 配置项按实际业务调整 GLUE_JOB_NAME = 'My_Job_Glue' TMP_FILE_SUFFIX = '.tmp' MANIFEST_SAVE_BUCKET = 'your-bucket-for-glue-config' MANIFEST_SAVE_PREFIX = 'pending-file-manifests/' def lambda_handler(event, context): bucket_file_map = {} # 遍历所有事件聚合有效文件 for record in event["Records"]: # 兼容SQS投递和S3直接触发两种消息格式 if "body" in record: s3_event_body = json.loads(record["body"]) s3_records = s3_event_body.get("Records", []) else: s3_records = [record] for s3_record in s3_records: bucket = s3_record["s3"]["bucket"]["name"] raw_key = s3_record["s3"]["object"]["key"] file_key = unquote_plus(raw_key) # 过滤临时写入文件 if file_key.endswith(TMP_FILE_SUFFIX): continue if bucket not in bucket_file_map: bucket_file_map[bucket] = set() bucket_file_map[bucket].add(file_key) # 无有效文件直接返回 if not bucket_file_map: return {"statusCode": 200, "body": json.dumps("No valid files to process")} # 检查Glue作业是否有运行中实例 job_run_list = glue.get_job_runs(JobName=GLUE_JOB_NAME, MaxResults=10) has_running_job = any( run["JobRunState"] in ["STARTING", "RUNNING", "WAITING", "STOPPING"] for run in job_run_list.get("JobRuns", []) ) # 生成待处理文件清单 file_path_list = [] for bucket, keys in bucket_file_map.items(): for key in keys: file_path_list.append(f"s3://{bucket}/{key}") # 清单存S3,避免Lambda传参过长 manifest_key = f"{MANIFEST_SAVE_PREFIX}{datetime.utcnow().strftime('%Y%m%d%H%M%S')}.txt" s3.put_object( Bucket=MANIFEST_SAVE_BUCKET, Key=manifest_key, Body="\n".join(file_path_list).encode("utf-8") ) # 有运行中作业直接返回,新文件已经存到待处理清单,下次作业启动统一处理 if has_running_job: return {"statusCode": 200, "body": json.dumps("Glue job is running, new files added to pending list")} # 无运行中作业,启动新的Glue作业 start_resp = glue.start_job_run( JobName=GLUE_JOB_NAME, Arguments={ "--manifest_bucket": MANIFEST_SAVE_BUCKET, "--manifest_prefix": MANIFEST_SAVE_PREFIX } ) return { "statusCode": 200, "body": json.dumps(f"Started Glue job {start_resp['JobRunId']}, total pending files: {len(file_path_list)}") }
配套调整
- Glue作业脚本需要改造读取逻辑:读取传入的
manifest_prefix路径下所有未处理的清单文件,解析出所有待处理的S3文件路径做加工,处理完成后删除对应清单文件,避免重复处理。 - 如果不想引入SQS,也可以用DynamoDB存储文件事件+TTL定时触发逻辑,但SQS方案运维成本更低,稳定性更高。
内容的提问来源于stack exchange,提问作者Rafael Souza
相关产品推荐
相关产品推荐

