如何实现带重试逻辑的AWS Lambda定时调度:基于S3文件状态触发动作
解决方案:AWS Lambda检查S3 _SUCCESS文件的调度优化
核心问题
之前的逻辑缺少状态跟踪,EventBridge固定时间触发Lambda时,不管该日期的任务是否已经完成(触发过DAG或已失败),都会重复执行,导致重复发送Slack通知或重复触发DAG。
解决思路
引入DynamoDB存储任务状态,记录每个日期的重试次数和任务状态,让Lambda每次执行前先判断当前任务的状态,避免重复操作;同时严格控制重试次数和通知时机。
具体实现步骤
1. 创建DynamoDB状态表
创建一张DynamoDB表,用于跟踪每个日期的任务状态:
- 主键:
dt(字符串类型,对应任务的日期格式YYYYMMDD) - 字段:
retry_count:数字类型,记录已重试次数status:字符串类型,可选值pending(待重试)、success(已成功触发DAG)、failed(已失败)
2. 修改Lambda函数逻辑
更新后的代码会先读取DynamoDB的状态,根据状态决定后续操作:
import boto3 from datetime import date, timedelta import requests import json import logging import botocore import ast import base64 from boto3.dynamodb.conditions import Key # 初始化客户端 s3 = boto3.resource('s3') dynamodb = boto3.resource('dynamodb') status_table = dynamodb.Table('YOUR_DYNAMODB_TABLE_NAME') # 替换为你的DynamoDB表名 # 配置参数 SLACK_WEBHOOK_URL = "" # 替换为你的Slack Webhook地址 MWAA_ENV_NAME = 'airflow-prod-env' DAG_NAME = 'sample_dag' MAX_RETRIES = 40 # 10小时/15分钟=40次重试 CHECK_DATE_DELTA = 2 # 检查两天前的日期 def send_slack_message(message): slack_payload = {'text': message} response = requests.post(SLACK_WEBHOOK_URL, json.dumps(slack_payload)) logging.info(f'Slack通知响应: {response.text}') def get_task_status(dt): """从DynamoDB获取指定日期的任务状态""" try: response = status_table.query(KeyConditionExpression=Key('dt').eq(dt)) if response['Items']: return response['Items'][0] else: # 首次执行,初始化状态 status_table.put_item(Item={ 'dt': dt, 'retry_count': 0, 'status': 'pending' }) return {'retry_count': 0, 'status': 'pending'} except Exception as e: logging.error(f'读取DynamoDB状态失败: {str(e)}') raise def update_task_status(dt, retry_count, status): """更新DynamoDB中的任务状态""" try: status_table.put_item(Item={ 'dt': dt, 'retry_count': retry_count, 'status': status }) except Exception as e: logging.error(f'更新DynamoDB状态失败: {str(e)}') raise def trigger_mwaa_dag(): """触发MWAA DAG运行""" try: client = boto3.client('mwaa') mwaa_cli_token = client.create_cli_token(Name=MWAA_ENV_NAME) conn = http.client.HTTPSConnection(mwaa_cli_token['WebServerHostname']) payload = f'dags trigger {DAG_NAME}' headers = { 'Authorization': f'Bearer {mwaa_cli_token["CliToken"]}', 'Content-Type': 'text/plain' } conn.request("POST", "/aws_mwaa/cli", payload, headers) res = conn.getresponse() data = res.read().decode("UTF-8") mydata = ast.literal_eval(data) result = base64.b64decode(mydata['stdout']).decode('utf-8') logging.info(f'DAG触发结果: {result}') return result except Exception as e: logging.error(f'触发DAG失败: {str(e)}') raise def check_success_file(bucket, path): """检查S3路径下是否存在_SUCCESS文件""" try: s3.Object(bucket, path).load() return True except botocore.exceptions.ClientError as e: if e.response['Error']['Code'] == '404': return False else: logging.error(f'S3访问错误: {str(e)}') raise def lambda_handler(event, context): # 初始化日期和路径 check_date = date.today() - timedelta(days=CHECK_DATE_DELTA) dt = check_date.strftime('%Y%m%d') bucket = "" # 替换为你的S3桶名 success_path = f"test_lambda/US/{dt}/_SUCCESS" # 获取当前任务状态 task_status = get_task_status(dt) current_status = task_status['status'] current_retry = task_status['retry_count'] # 任务已完成(成功或失败),直接终止 if current_status in ['success', 'failed']: logging.info(f'日期{dt}的任务已处于{current_status}状态,无需执行') return # 检查_SUCCESS文件 file_exists = check_success_file(bucket, success_path) if file_exists: # 文件存在,触发DAG并标记任务成功 logging.info(f'找到_SUCCESS文件,路径: {bucket}/{success_path}') trigger_mwaa_dag() update_task_status(dt, current_retry, 'success') logging.info(f'日期{dt}的任务已标记为success,终止调度') else: # 文件不存在,处理重试逻辑 new_retry_count = current_retry + 1 if new_retry_count >= MAX_RETRIES: # 达到最大重试次数,发送失败通知并标记任务失败 error_msg = f'日期{dt}的_SUCCESS文件未找到,已重试{MAX_RETRIES}次,任务失败' logging.error(error_msg) send_slack_message(error_msg) update_task_status(dt, new_retry_count, 'failed') logging.info(f'日期{dt}的任务已标记为failed') else: # 未达到最大重试次数,更新重试计数,不发通知 update_task_status(dt, new_retry_count, 'pending') logging.info(f'日期{dt}的_SUCCESS文件未找到,已重试{new_retry_count}次,等待下一次执行')
3. 调整EventBridge调度
保持原有的Cron表达式(0/15 4-13 ? * MON-FRI *)即可,因为Lambda会通过DynamoDB状态自动跳过已完成的任务,不会重复执行操作。
内容的提问来源于stack exchange,提问作者Dirie
相关产品推荐
相关产品推荐

