如何配置Airflow DAG按月特定日期触发且仅在文件更新时执行
我来给你梳理下这个需求的具体实现方案,结合Airflow的特性,有两种比较靠谱的方式,还有一些优化细节可以参考:
核心思路
你的需求是每月特定日期触发DAG,但仅当目录内表文件有更新时才执行迁移。Airflow本身的调度是基于时间触发的,所以我们需要把「时间调度」和「文件更新检查」结合起来——要么在调度触发后先做检查,只有通过才执行迁移;要么让DAG在指定日期窗口内监控文件更新,有更新才触发迁移。
方案一:固定日期调度+前置检查任务(推荐)
这种方式逻辑最简单,适合每月固定日期执行一次检查,有更新就跑迁移,没有就直接结束。
实现步骤
- 定义Cron调度规则:设置
schedule_interval为你需要的每月特定日期,比如每月15号凌晨0点:'0 0 15 * *'。 - 添加文件更新检查任务:用
PythonOperator写一个检查函数,对比文件的内容哈希(比修改时间更精准),判断是否有更新。 - 控制迁移任务的执行:只有检查任务成功,才执行后续的迁移任务。
代码示例
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结合分支任务实现。
实现步骤
- 先判断是否处于指定日期:用
PythonOperator检查当前执行日期是否是目标日期。 - 分支任务控制流程:如果是指定日期,启动
FileSensor监控文件更新;否则直接结束DAG。 - 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资源,适合全天监控。
优化建议
- 避免Variable冲突:如果有多个类似DAG,建议用XCom或专门的元数据表存储哈希/修改时间,而不是全局Variable。
- 过滤临时文件:在检查文件时,排除
.tmp、.part等临时文件,避免误触发。 - 日志记录:在检查函数中添加详细日志,方便排查问题,比如记录当前哈希值、上次哈希值等。
- 异常处理:针对目录不存在、文件权限不足等情况添加异常捕获,让任务失败原因更清晰。
内容的提问来源于stack exchange,提问作者rajat_th
相关产品推荐
相关产品推荐

