如何配置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.comAIRFLOW_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
相关产品推荐
相关产品推荐

