Airflow设置@daily调度 DAG解析时每30秒执行代码而非每日触发
问题根因
你的Mongo插入逻辑写在了DAG文件的顶层作用域,没有被封装为Airflow可调度的任务单元,所以每次DAG解析流程加载文件时,这段代码都会直接执行,和schedule_interval配置、min_file_process_interval参数没有关系。
你代码里的具体问题:
- 直接在DAG定义上下文里实例化
MongoHook后立刻调用insert_one()方法,这段代码不属于任何任务的执行逻辑,属于DAG文件加载时就会运行的顶层代码,解析一次跑一次。 MongoHook本身不是Operator,只是封装Mongo连接逻辑的工具类,不能直接作为任务挂载到DAG中。- 修改webserver容器内的
airflow.cfg本身就不会生效:DAG解析、任务调度的逻辑是scheduler组件负责的,webserver只负责页面展示,改webserver的配置影响不到解析和调度流程。
修复方法
把Mongo插入逻辑封装为标准Airflow任务即可,两种常用实现方式:
方式1:用PythonOperator包装Hook调用
适合需要在插入前后加自定义处理逻辑的场景,修复后的完整可运行代码:
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.mongo.hooks.mongo import MongoHook default_args = { 'owner': 'kw', 'retries': 1, 'retry_delay': timedelta(minutes=15), } doc_raw = { 'name': "Chaitanya" } MONGO_CONN_ID = 'mongo_conn' MONGO_COLLECTION = 'airflowtest' MONGO_DB = 'airflowtest' def print_hello(): return 'Hello world from first Airflow DAG!' def mongo_insert_exec(): # Hook实例化和方法调用放在任务函数内部,只有任务被调度触发时才会执行 mongo_hook = MongoHook(mongo_conn_id=MONGO_CONN_ID) mongo_hook.insert_one( mongo_collection=MONGO_COLLECTION, doc=doc_raw, mongo_db=MONGO_DB ) with DAG( dag_id='kw_mongo_test4', default_args=default_args, start_date=datetime(2022, 6, 27), schedule_interval="@daily", description='use case of mongo operator in airflow', catchup=False ) as dag: task1 = PythonOperator( task_id='hello_task', python_callable=print_hello ) task2 = PythonOperator( task_id="mongo_insert_test", python_callable=mongo_insert_exec ) # 配置任务依赖 task1 >> task2
方式2:直接使用官方MongoInsertOperator
不需要手动封装Hook,代码更简洁,使用前先确认已经安装对应依赖:pip install apache-airflow-providers-mongo
对应task2的写法替换为:
from airflow.providers.mongo.operators.mongo import MongoInsertOperator task2 = MongoInsertOperator( task_id="mongo_insert_test", mongo_conn_id=MONGO_CONN_ID, mongo_collection=MONGO_COLLECTION, mongo_db=MONGO_DB, doc=doc_raw )
注意事项
- Airflow的DAG解析本质是周期性执行所有DAG文件的顶层代码,以此扫描DAG结构、生成任务依赖拓扑,所有写在顶层作用域的业务代码都会在每次解析时直接运行,所有业务逻辑必须封装在任务的可调用函数、或者Operator的执行逻辑内部,绝对不能直接写在顶层。
- 涉及DAG解析、调度频率的配置项,需要修改scheduler服务对应的配置文件,仅修改webserver的配置不会对解析、调度逻辑产生任何影响。
内容的提问来源于stack exchange,提问作者Krystian Warda
相关产品推荐
相关产品推荐

