如何用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
- 建议指定Airflow具体版本(如
- 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
- DAG文件放入
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
相关产品推荐
相关产品推荐

