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

如何用API触发DAG?Airflow-Docker项目集成Flask及步骤排查

现有操作错误检查
  • Dockerfile问题
    • 建议指定Airflow具体版本(如apache/airflow:2.8.0),避免latest版本不稳定
    • 复制requirements.txt到/opt/airflow/目录(而非根目录),避免airflow用户无写入权限,修改后Dockerfile如下:
      FROM apache/airflow:2.8.0
      
      USER airflow
      
      COPY requirements.txt /opt/airflow/
      
      RUN pip install --no-cache-dir "apache-airflow==${AIRFLOW_VERSION}" -r /opt/airflow/requirements.txt
      
  • docker-compose.yml问题
    • 构建镜像时需显式打标签:docker build -t pythonairflow:latest .,否则compose会拉取不存在的公共镜像
    • 添加环境变量关闭示例DAG、设置时区:
      environment:
        - AIRFLOW__CORE__LOAD_EXAMPLES=false
        - TZ=Asia/Shanghai
      
  • DAG文件问题
    • 未持久化训练好的模型与标准化器,导致后续Flask无法调用预测
    • 导入requests但未使用,可删除
  • 操作流程问题
    • DAG文件放入airflow/dags后无需重启容器,Airflow默认每30秒自动扫描DAG目录
    • 第一次启动standalone模式后,需保存airflow/standalone_admin_password.txt中的初始密码,否则无法登录Web UI
Flask集成步骤

1. 更新依赖

修改requirements.txt添加Flask与模型持久化依赖:

scikit-learn
flask
joblib

2. 修改DAG的训练逻辑,添加模型保存

在train函数末尾添加模型与标准化器的持久化代码:

import os
import joblib

# 创建模型保存目录并确保权限
model_dir = '/opt/airflow/models'
os.makedirs(model_dir, exist_ok=True, mode=0o755)

# 保存模型与标准化器
joblib.dump(model, os.path.join(model_dir, 'linear_regression_model.pkl'))
joblib.dump(scaler, os.path.join(model_dir, 'scaler.pkl'))

3. 创建Flask服务文件

在项目根目录新建flask_app.py:

from flask import Flask, request, jsonify
import joblib
import numpy as np

app = Flask(__name__)

# 加载模型与标准化器
model_path = '/opt/airflow/models/linear_regression_model.pkl'
scaler_path = '/opt/airflow/models/scaler.pkl'
model, scaler = None, None

try:
    model = joblib.load(model_path)
    scaler = joblib.load(scaler_path)
except FileNotFoundError:
    print("请先运行DAG训练模型")

@app.route('/predict', methods=['POST'])
def predict():
    if not model or not scaler:
        return jsonify({'error': '模型未训练,请先执行DAG'}), 500
    
    data = request.get_json()
    if not data or 'features' not in data:
        return jsonify({'error': '请求缺少features字段'}), 400
    
    try:
        features = np.array(data['features']).reshape(1, -1)
        scaled_features = scaler.transform(features)
        prediction = model.predict(scaled_features)[0]
        return jsonify({'prediction': float(prediction)})
    except Exception as e:
        return jsonify({'error': str(e)}), 400

if __name__ == '__main__':
    app.run(host='0.0.0.0', port=5000, debug=False)

4. 更新容器配置

  • 修改docker-compose.yml,添加Flask端口映射、挂载Flask文件、调整启动命令:
    version: '3'
    
    services:
      sleek-airflow:
        image: pythonairflow:latest
        volumes:
          - ./airflow:/opt/airflow
          - ./flask_app.py:/opt/airflow/flask_app.py
        ports:
          - "8080:8080"
          - "5000:5000"
        environment:
          - AIRFLOW__CORE__LOAD_EXAMPLES=false
          - TZ=Asia/Shanghai
        command: >
          bash -c "airflow standalone & python /opt/airflow/flask_app.py"
    

5. 测试接口

运行DAG训练模型后,用curl发起预测请求:

curl -X POST -H "Content-Type: application/json" -d '{"features": [8.3252, 41.0, 6.984126984, 1.023809524, 322.0, 2.555555556, 37.88, -122.23]}' http://localhost:5000/predict
通过API触发DAG的方法

1. 获取API Token

  • 登录Airflow Web UI(http://localhost:8080)
  • 点击右上角头像 → Profile → API Keys → Create API Key,保存生成的Token

2. 触发DAG

用curl调用Airflow REST API触发DAG:

curl -X POST \
  http://localhost:8080/api/v1/dags/pipeline_dag/dagRuns \
  -H "Authorization: Bearer 你的API_TOKEN" \
  -H "Content-Type: application/json" \
  -d '{"conf": {}}'

如需传递参数,可在conf中添加键值对,例如{"conf": {"param": "value"}},在DAG函数中通过context['dag_run'].conf.get('param')获取。

3. 验证触发结果

  • 查看Airflow Web UI的DAG Runs页面,确认新的运行记录
  • 或用API查询运行状态:
    curl -H "Authorization: Bearer 你的API_TOKEN" http://localhost:8080/api/v1/dags/pipeline_dag/dagRuns
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 03:54:54