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

如何优化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作业处理全量文件。

具体实现步骤

    1. 调整事件链路:S3新对象创建事件不要直接投递Lambda,先投递到SQS队列做缓冲。队列参数配置:
    • 消息可见性超时 = Kafka Connect最大空闲刷写时间 + 2分钟(比如Kafka空闲刷写阈值是5分钟,就设为7分钟)
    • 消息保留期设为1天即可
    1. 配置Lambda触发源为上述SQS队列:
    • 批处理大小设为20(覆盖单批次Kafka写入的最大文件拆分数量即可)
    • 批处理窗口设为60秒,攒齐窗口内的所有事件再触发Lambda
    1. 改造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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 14:01:14