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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 06:42:19