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

能否在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:05:12