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

如何配置Airflow DAG按月特定日期触发且仅在文件更新时执行

我来给你梳理下这个需求的具体实现方案,结合Airflow的特性,有两种比较靠谱的方式,还有一些优化细节可以参考:

核心思路

你的需求是每月特定日期触发DAG,但仅当目录内表文件有更新时才执行迁移。Airflow本身的调度是基于时间触发的,所以我们需要把「时间调度」和「文件更新检查」结合起来——要么在调度触发后先做检查,只有通过才执行迁移;要么让DAG在指定日期窗口内监控文件更新,有更新才触发迁移。

方案一:固定日期调度+前置检查任务(推荐)

这种方式逻辑最简单,适合每月固定日期执行一次检查,有更新就跑迁移,没有就直接结束。

实现步骤

  1. 定义Cron调度规则:设置schedule_interval为你需要的每月特定日期,比如每月15号凌晨0点:'0 0 15 * *'。
  2. 添加文件更新检查任务:用PythonOperator写一个检查函数,对比文件的内容哈希(比修改时间更精准),判断是否有更新。
  3. 控制迁移任务的执行:只有检查任务成功,才执行后续的迁移任务。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.models import Variable
from datetime import datetime
from pathlib import Path
import hashlib

# 计算目录下所有表文件的哈希值(精准判断内容是否更新)
def get_table_files_hash(target_dir):
    hash_obj = hashlib.md5()
    # 按文件名排序,保证每次哈希计算的文件顺序一致
    for file_path in sorted(Path(target_dir).glob("*.csv")):  # 替换成你的文件后缀
        if file_path.is_file():
            with open(file_path, 'rb') as f:
                while chunk := f.read(4096):
                    hash_obj.update(chunk)
    return hash_obj.hexdigest()

# 检查文件是否更新的函数
def check_file_update(**context):
    target_dir = "/path/to/your/table_files"  # 替换成你的目标目录
    current_hash = get_table_files_hash(target_dir)
    
    # 从Airflow变量中获取上次记录的哈希值,首次执行默认空字符串
    last_recorded_hash = Variable.get("last_table_files_hash", default_var="")
    
    if current_hash != last_recorded_hash:
        # 更新变量为最新哈希值
        Variable.set("last_table_files_hash", current_hash)
        print("Detected updates in table files, proceeding with migration.")
        return True
    else:
        # 没有更新,抛出异常让任务失败,后续迁移任务不会执行
        raise ValueError("No content updates found in table files directory.")

# 你的迁移逻辑函数(替换成实际的迁移代码)
def execute_table_migration(**context):
    print("Starting table migration...")
    # 这里写你的迁移逻辑,比如读取文件写入数据库等
    # ...
    print("Table migration completed successfully.")

# DAG配置
default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 0,  # 无更新时无需重试
}

with DAG(
    'monthly_table_migration',
    default_args=default_args,
    schedule_interval='0 0 15 * *',  # 每月15号凌晨0点触发
    catchup=False,  # 不补跑历史任务
    tags=['table_migration']
) as dag:

    check_update = PythonOperator(
        task_id='check_table_file_updates',
        python_callable=check_file_update,
        provide_context=True
    )

    migrate_task = PythonOperator(
        task_id='run_table_migration',
        python_callable=execute_table_migration,
        provide_context=True
    )

    # 任务依赖:检查通过才执行迁移
    check_update >> migrate_task

方案说明

  • 用哈希值对比代替修改时间,避免因touch等操作误判文件更新,更精准。
  • 用Airflow的Variable存储上次的哈希值,全局可访问,下次执行时直接对比。
  • 如果检查任务失败,整个DAG的迁移环节会被跳过,完全符合需求。
方案二:指定日期窗口内监控文件更新

如果需要在每月特定日期的一整天内,随时监控文件更新(比如文件可能在15号的任意时间上传),可以用FileSensor结合分支任务实现。

实现步骤

  1. 先判断是否处于指定日期:用PythonOperator检查当前执行日期是否是目标日期。
  2. 分支任务控制流程:如果是指定日期,启动FileSensor监控文件更新;否则直接结束DAG。
  3. Sensor触发迁移:当Sensor检测到文件更新时,执行迁移任务。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.branch import BranchPythonOperator
from airflow.sensors.filesystem import FileSensor
from datetime import datetime, timedelta

# 检查是否是指定日期(比如每月15号)
def is_target_date(**context):
    execution_date = context['execution_date']
    return execution_date.day == 15

# 空任务:非指定日期时直接结束
def end_dag(**context):
    print("Not the scheduled date, exiting DAG.")

# DAG配置
default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5)
}

with DAG(
    'table_migration_sensor_based',
    default_args=default_args,
    schedule_interval='@daily',  # 每天触发一次
    catchup=False,
    tags=['table_migration']
) as dag:

    check_date = PythonOperator(
        task_id='check_if_target_date',
        python_callable=is_target_date,
        provide_context=True,
        do_xcom_push=True
    )

    branch_task = BranchPythonOperator(
        task_id='branch_based_on_date',
        python_callable=lambda ctx: 'wait_for_file_updates' if ctx['ti'].xcom_pull(task_ids='check_if_target_date') else 'end_dag_task'
    )

    end_dag_task = PythonOperator(
        task_id='end_dag_task',
        python_callable=end_dag,
        provide_context=True
    )

    # 监控目标目录下的文件更新,每小时检查一次,超时时间为1天(仅在指定日期内监控)
    file_sensor = FileSensor(
        task_id='wait_for_file_updates',
        filepath='/path/to/your/table_files/*.csv',  # 匹配你的文件
        poke_interval=3600,  # 每小时检查一次
        mode='reschedule',  # 非检查时段释放worker资源
        timeout=86400,  # 1天后超时
        on_failure_callback=lambda ctx: print("No file updates detected on scheduled date.")
    )

    migrate_task = PythonOperator(
        task_id='run_table_migration',
        python_callable=execute_table_migration,  # 同方案一的迁移函数
        provide_context=True
    )

    # 任务依赖
    check_date >> branch_task >> [file_sensor, end_dag_task]
    file_sensor >> migrate_task

方案说明

  • 用@daily调度每天触发,然后通过分支任务过滤出指定日期的执行。
  • FileSensor的mode='reschedule'可以避免长时间占用worker资源,适合全天监控。
优化建议
  1. 避免Variable冲突:如果有多个类似DAG,建议用XCom或专门的元数据表存储哈希/修改时间,而不是全局Variable。
  2. 过滤临时文件:在检查文件时,排除.tmp、.part等临时文件,避免误触发。
  3. 日志记录:在检查函数中添加详细日志,方便排查问题,比如记录当前哈希值、上次哈希值等。
  4. 异常处理:针对目录不存在、文件权限不足等情况添加异常捕获,让任务失败原因更清晰。

内容的提问来源于stack exchange,提问作者rajat_th

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:05:26