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

如何实现带重试逻辑的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 13:20:48