能否在Airflow Web Server UI添加按钮上传文件并触发DAG解析?
需求可行性分析与实现方案
这个需求完全可行,Airflow支持通过自定义扩展实现带文件上传的DAG触发功能,效果可以对标「Trigger DAG w/ config」。以下是具体实现思路:
扩展Airflow UI添加上传按钮
借助Airflow的插件系统自定义Web UI组件,创建一个包含文件选择器的独立页面或嵌入现有DAG列表页的按钮。通过Flask蓝图搭建后台接口,前端用表单实现文件选择与提交逻辑,点击按钮即可弹出文件选择器。文件上传后的存储处理
后台接收文件后,将其存储到Airflow集群可统一访问的位置:比如共享本地目录(需确保所有Worker节点挂载该目录)、S3/GCS等分布式存储。存储完成后记录文件的完整路径或唯一标识,作为DAG运行的配置参数。关联DAG触发逻辑
文件存储完成后,调用Airflow的内部API(DagRun.create方法)触发目标DAG,并将文件路径传入DAG的config中。DAG启动后,就能通过context['dag_run'].conf获取文件路径,进而执行解析逻辑。核心代码示例
自定义上传插件的后台逻辑:from airflow.plugins_manager import AirflowPlugin from flask import Blueprint, request, redirect, url_for from airflow.models import DagRun from airflow.utils.state import State import os upload_bp = Blueprint( 'log_upload_bp', __name__, template_folder='templates' ) @upload_bp.route('/upload-log', methods=['GET', 'POST']) def upload_and_trigger(): if request.method == 'POST': log_file = request.files['log_file'] # 确保上传目录存在且权限正确 upload_dir = '/airflow/shared/log_uploads' os.makedirs(upload_dir, exist_ok=True) file_path = os.path.join(upload_dir, log_file.filename) log_file.save(file_path) # 触发目标DAG并传入文件路径 DagRun.create( dag_id='log_parser_dag', conf={'log_path': file_path}, state=State.RUNNING ) return redirect(url_for('Airflow.index')) # 前端上传表单 return ''' <div style="padding: 20px;"> <h3>上传日志并触发解析DAG</h3> <form method="post" enctype="multipart/form-data"> <input type="file" name="log_file" accept=".log,.txt" required> <button type="submit" style="margin-left: 10px;">上传并运行</button> </form> </div> ''' class LogUploadPlugin(AirflowPlugin): name = 'log_upload_plugin' flask_blueprints = [upload_bp]DAG中的文件解析逻辑:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def parse_log(**context): log_path = context['dag_run'].conf.get('log_path') if not log_path or not os.path.exists(log_path): raise ValueError("日志文件路径无效或不存在") # 这里替换为实际的日志解析逻辑 with open(log_path, 'r', encoding='utf-8') as f: line_count = len(f.readlines()) print(f"日志文件共{line_count}行,开始执行解析...") with DAG( dag_id='log_parser_dag', schedule_interval=None, start_date=datetime(2024, 1, 1), catchup=False ) as dag: parse_task = PythonOperator( task_id='parse_log_file', python_callable=parse_log, provide_context=True )关键注意事项
- 确保上传目录的读写权限对Airflow Webserver和Worker进程开放。
- 大文件优先使用分布式存储,避免本地存储的节点兼容性问题。
- 可添加文件类型校验、大小限制等安全逻辑,防止恶意文件上传。
- Airflow 2.x版本的插件写法与1.x略有差异,需根据实际版本调整代码。
内容的提问来源于stack exchange,提问作者Prinz Piuz
相关产品推荐
相关产品推荐

