如何在Airflow任务中访问通过REST API传入的JSON数据?
解决Airflow REST API触发DAG时自定义配置无法访问的问题
问题根源
你的请求代码存在两个关键错误:
- 自定义配置未正确放入DAG Run请求体的
conf字段,而是单独通过json参数传递,导致Airflow无法识别 - 同时使用
data和json参数发送请求,requests会优先采用json参数的内容,使得原本的dag_run_id等参数未被正确提交
修正后的请求代码
import requests from requests.auth import HTTPBasicAuth # 自定义配置数据 custom_conf = { "flag": "flag", "files": "files", "upload_path": "config.UPLOAD_FOLDER", "tmp_path": "config.TMP_FOLDER", "dataset_id": "dataset_id", "dicom_meta_data": "dicom_meta_data", "user_id": "request.user.id", "protocol": "http if request.is_secure() else http", "current_site": "request.get_host()", "deidentify": "deidentify", "email": "request.user.email", } # 构造完整的DAG Run请求体,将自定义配置放入conf字段 payload = { "conf": custom_conf, "dag_run_id": "trigger_16", "logical_date": "2022-08-29T11:33:49.726Z", } headers = { 'Content-type': 'application/json', 'Accept': 'application/json' } # 使用requests的json参数直接传递payload,无需手动序列化JSON r = requests.post( "http://localhost:8080/api/v1/dags/dag_1/dagRuns", auth=HTTPBasicAuth("airflow", "airflow"), json=payload, headers=headers ) print(r.status_code) print(r.text)
在DAG任务中访问自定义配置
在首个任务中,需通过Airflow的上下文context获取dag_run对象,进而读取传递的配置:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def first_task(**context): # 从上下文提取dag_run的配置 custom_config = context['dag_run'].conf # 读取具体配置字段 flag = custom_config.get('flag') user_id = custom_config.get('user_id') print(f"获取到的flag: {flag}, user_id: {user_id}") with DAG( dag_id='dag_1', start_date=datetime(2022, 8, 29), catchup=False ) as dag: first_task = PythonOperator( task_id='first_task', python_callable=first_task, provide_context=True # 必须开启该参数以传递上下文 )
注意事项
- 确保使用Airflow 2.x版本,代码中的API路径
/api/v1/dags/...仅适用于2.x版本 dag_run_id需保证全局唯一,重复的ID会触发Airflow的重复运行错误- 若使用模板字符串(如Jinja2),可直接通过
{{ dag_run.conf.flag }}引用配置字段
内容的提问来源于stack exchange,提问作者NeedToCodeAgain
相关产品推荐
相关产品推荐

