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

如何配置AWS Lambda触发Apache Airflow DAG处理S3文件入库MySQL

别担心,这个方案其实很成熟,我来一步步带你配置Lambda触发Airflow DAG,全程都是实操步骤,跟着来就行~

配置Lambda触发Airflow DAG的完整步骤

1. 先搞定Airflow的API访问权限

首先,Airflow 2.x自带稳定的REST API,这是我们触发DAG的核心入口,先做好基础配置:

  • 登录Airflow UI,进入Admin > Users,可以专门创建一个服务用户(避免用个人账号),点击编辑后选择Generate API Token,复制生成的令牌保存好,后面会用到
  • 确认Airflow Webserver的可达性:如果Airflow部署在公网,确保8080端口对外开放;如果在VPC内,需要把Lambda配置到同一个VPC的子网和安全组,且安全组允许Lambda访问Airflow的端口

2. 创建Lambda执行角色(IAM)

Lambda需要足够的权限来调用Airflow API、处理S3事件,还要能输出日志用于调试:

  • 登录AWS IAM控制台,创建新角色,类型选Lambda
  • 附加以下权限策略:
    • AWSLambdaBasicExecutionRole:给Lambda日志输出权限,方便排查问题
    • 如果Airflow在VPC内,额外附加AmazonVPCFullAccess(或者更精细的权限,比如允许访问指定子网、安全组)
  • 保存这个角色,后面创建Lambda函数时会用到

3. 编写Lambda触发Airflow的代码

我们用Python编写Lambda函数,核心是调用Airflow的dagRuns API来触发指定DAG,同时把S3上传的文件路径传给Airflow:

  • 登录AWS Lambda控制台,创建新函数,选择Author from scratch,运行时选Python 3.10+,执行角色选刚才创建的IAM角色
  • 在函数代码编辑器里替换成下面的代码:
import requests
import os

# 从环境变量读取配置,避免硬编码敏感信息
AIRFLOW_API_URL = os.environ["AIRFLOW_API_URL"]
AIRFLOW_DAG_ID = os.environ["AIRFLOW_DAG_ID"]
AIRFLOW_API_TOKEN = os.environ["AIRFLOW_API_TOKEN"]

def lambda_handler(event, context):
    # 从S3事件中提取所有上传的CSV文件路径
    s3_records = event["Records"]
    file_paths = [
        f"s3://{record['s3']['bucket']['name']}/{record['s3']['object']['key']}"
        for record in s3_records
    ]

    # 构建触发DAG的请求
    headers = {
        "Authorization": f"Bearer {AIRFLOW_API_TOKEN}",
        "Content-Type": "application/json"
    }
    payload = {
        "conf": {
            "s3_file_paths": file_paths  # 把文件路径传给Airflow DAG
        }
    }

    try:
        response = requests.post(
            f"{AIRFLOW_API_URL}/api/v1/dags/{AIRFLOW_DAG_ID}/dagRuns",
            headers=headers,
            json=payload
        )
        response.raise_for_status()  # 抛出HTTP错误便于调试
        print(f"✅ DAG {AIRFLOW_DAG_ID} 触发成功,待处理文件:{file_paths}")
        return {
            "statusCode": 200,
            "body": f"DAG触发成功,文件列表:{file_paths}"
        }
    except Exception as e:
        print(f"❌ 触发DAG失败:{str(e)}")
        return {
            "statusCode": 500,
            "body": f"触发错误:{str(e)}"
        }
  • 配置Lambda的环境变量:
    • AIRFLOW_API_URL:你的Airflow Webserver地址,比如http://airflow-webserver:8080(VPC内)或者公网域名https://your-airflow-domain.com
    • AIRFLOW_DAG_ID:你要触发的Airflow DAG的ID(和DAG代码里的dag_id值一致)
    • AIRFLOW_API_TOKEN:刚才在Airflow里生成的API令牌
  • 点击部署保存函数

4. 配置S3触发Lambda

现在让S3在有CSV文件上传到MY_BUCKET/MY_DIRECTORY时自动触发Lambda:

  • 登录AWS S3控制台,找到MY_BUCKET桶,进入属性标签页,拉到最下方的事件通知
  • 点击创建事件通知:
    • 通知名称:自定义一个,比如csv_upload_trigger
    • 前缀:填MY_DIRECTORY/(注意末尾的斜杠,确保只触发该文件夹下的文件)
    • 后缀:填.csv(只响应CSV文件上传)
    • 事件类型:勾选所有对象创建事件(或者只选Put事件,根据你的上传方式)
    • 目标:选Lambda函数,然后选择你刚才创建的Lambda函数
  • 保存配置

5. Airflow DAG适配(处理CSV并写入MySQL)

你的DAG需要能接收Lambda传递的文件路径,然后完成CSV解析和MySQL写入,这里给你一个极简示例:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.providers.mysql.hooks.mysql import MySqlHook
from datetime import datetime
import pandas as pd
from io import StringIO

def process_and_load_csv(**context):
    # 获取Lambda传递的文件路径列表
    s3_file_paths = context["dag_run"].conf.get("s3_file_paths", [])
    if not s3_file_paths:
        print("⚠️ 没有收到待处理的文件路径")
        return

    # 初始化S3和MySQL Hook(需提前在Airflow Connections里配置aws_default和mysql_default)
    s3_hook = S3Hook(aws_conn_id="aws_default")
    mysql_hook = MySqlHook(mysql_conn_id="mysql_default")

    for file_path in s3_file_paths:
        # 拆分bucket和key
        bucket = file_path.split("/")[2]
        key = "/".join(file_path.split("/")[3:])
        
        # 读取S3上的CSV文件
        csv_content = s3_hook.read_key(key=key, bucket_name=bucket)
        df = pd.read_csv(StringIO(csv_content))
        
        # 写入MySQL(假设你的表结构和CSV列匹配,需提前创建好目标表)
        mysql_hook.insert_rows(
            table="your_target_table",
            rows=df.values.tolist(),
            target_fields=df.columns.tolist()
        )
        print(f"✅ 处理完成:{file_path}")

with DAG(
    dag_id="s3_csv_to_mysql",  # 这个ID要和Lambda里的AIRFLOW_DAG_ID一致
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # 由外部触发,所以设为None
    catchup=False,
    tags=["s3", "mysql", "lambda"]
) as dag:
    process_task = PythonOperator(
        task_id="process_csv_files",
        python_callable=process_and_load_csv,
        provide_context=True
    )
  • 注意:需要在Airflow的Admin > Connections里配置好aws_default(AWS连接)和mysql_default(MySQL连接)的信息
  • 确保DAG在Airflow UI里处于Unpaused状态,否则即使触发也不会运行

一些避坑提示

  • 如果一次性上传1000个文件,S3会发送多个事件给Lambda,可能导致多次触发DAG。如果希望等所有文件上传完再触发一次,可以考虑用SQS做中间层:把S3事件发送到SQS,Lambda设置批量接收消息,或者加延迟逻辑收集一段时间内的所有文件后再触发DAG
  • 生产环境不要把Airflow Webserver暴露在公网,尽量放在VPC内,让Lambda在同一个VPC内访问,更安全
  • 给Lambda配置足够的内存和超时时间,避免处理大文件时超时
  • 可以在Airflow的DAG Runs页面查看触发记录,在Lambda的监控 > 日志里查看执行日志,方便调试

内容的提问来源于stack exchange,提问作者Joy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 14:22:48