如何通过AirFlow REST API带conf变量实现DAG每日定时调度?
可以实现,提供三种常用方案:
方案1:动态生成独立DAG(推荐)
直接基于端点列表循环生成多个配置好的DAG,每个DAG对应一个端点,自带每日调度:
- 定义端点列表(可存在AirFlow Variables或直接写在代码中)
- 循环每个端点,创建独立的DAG实例,设置
schedule_interval="@daily",并将对应端点信息写入DAG的默认conf或直接传入任务 - 爬取任务中通过
context['dag_run'].conf获取端点参数
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime # 端点列表,可替换为从AirFlow Variables读取 ENDPOINTS = ["python", "java", "ruby"] def crawl_website(**context): endpoint = context['dag_run'].conf.get('endpoint') # 这里写你的Selenium爬取逻辑,使用endpoint拼接目标URL print(f"爬取端点: {endpoint}") for endpoint in ENDPOINTS: dag_id = f"crawl_website_{endpoint}" default_args = { 'start_date': datetime(2024, 1, 1), 'catchup': False } with DAG( dag_id=dag_id, default_args=default_args, schedule_interval="@daily", user_defined_macros={'endpoint': endpoint} ) as dag: crawl_task = PythonOperator( task_id=f"crawl_{endpoint}", python_callable=crawl_website, provide_context=True )
这种方式每个端点任务完全独立,便于单独监控、暂停或重启,失败后不会影响其他端点的爬取。
方案2:用TriggerDagRunOperator批量触发原DAG
复用你已有的爬取DAG,新建一个调度DAG,每日运行并触发原DAG多次,每次传递不同的conf:
- 新建一个调度DAG,设置
schedule_interval="@daily" - 循环端点列表,为每个端点创建一个
TriggerDagRunOperator,指定触发目标DAG,并传入对应端点的conf - 原爬取DAG无需设置调度,仅通过被触发执行
示例代码:
from airflow import DAG from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime ENDPOINTS = ["python", "java", "ruby"] TARGET_DAG_ID = "your_existing_crawl_dag" default_args = { 'start_date': datetime(2024, 1, 1), 'catchup': False } with DAG( dag_id="trigger_crawl_dags_daily", default_args=default_args, schedule_interval="@daily" ) as dag: for endpoint in ENDPOINTS: trigger_task = TriggerDagRunOperator( task_id=f"trigger_crawl_{endpoint}", trigger_dag_id=TARGET_DAG_ID, conf={'endpoint': endpoint}, wait_for_completion=False # 根据需求选择是否等待任务完成 )
这种方式无需修改原DAG,仅新增一个调度控制DAG,适合需要统一管理触发逻辑的场景。
方案3:同一DAG内生成多任务
如果不需要任务完全独立,可在同一个DAG中通过TaskGroup循环生成多个爬取任务,每日统一调度:
- 在原DAG中使用
TaskGroup,循环端点列表创建爬取任务 - 每个任务直接接收端点参数,无需通过conf传递
示例代码:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup from datetime import datetime ENDPOINTS = ["python", "java", "ruby"] def crawl_website(endpoint): # 这里写你的Selenium爬取逻辑 print(f"爬取端点: {endpoint}") default_args = { 'start_date': datetime(2024, 1, 1), 'catchup': False } with DAG( dag_id="crawl_all_endpoints_daily", default_args=default_args, schedule_interval="@daily" ) as dag: with TaskGroup("crawl_tasks") as crawl_group: for endpoint in ENDPOINTS: crawl_task = PythonOperator( task_id=f"crawl_{endpoint}", python_callable=crawl_website, op_kwargs={'endpoint': endpoint} )
这种方式所有任务在同一个DAG下,便于查看整体执行状态,适合端点数量较少、无需单独管控的场景。
内容的提问来源于stack exchange,提问作者Jooo21
相关产品推荐
相关产品推荐

