Airflow周/月调度:以处理日期为结束日期时如何选起始日期
解决Airflow周调度DAG处理对应周期文件的问题
方案1:结合data_interval_end宏与手动触发参数
这是最适配你需求的方案:自动调度时用data_interval_end匹配文件到达日期,文件延迟时通过手动传入参数指定目标日期。
实现步骤:
- 任务函数优先读取手动触发传入的配置参数,无参数时使用周期结束日期:
from airflow.decorators import task from airflow.models import DAG from datetime import datetime def process_s3_file(**context): # 读取手动触发时传入的目标日期 dag_run_conf = context.get('dag_run').conf or {} target_date = dag_run_conf.get('target_date') if not target_date: # 自动调度时取周期结束日期(即文件到达的周一) target_date = context['data_interval_end'].strftime('%Y-%m-%d') # 按target_date匹配S3文件并处理 print(f"Processing file for date: {target_date}") # 你的文件处理逻辑 with DAG( dag_id='weekly_s3_file_processor', start_date=datetime(2023, 11, 6), schedule_interval='30 6 * * 1', # 每周一6:30 CST触发 catchup=False, ) as dag: process_task = task.python( task_id='process_s3_file', python_callable=process_s3_file, provide_context=True )
- 文件延迟时手动触发:
在Airflow UI触发DAG时,填写Conf为{"target_date": "2023-11-13"},任务会直接处理指定日期的文件。
方案2:日期偏移宏适配自动调度
如果不想依赖data_interval_end,可对ds宏做周偏移,同时保留手动触发的灵活性:
实现步骤:
任务中用macros.ds_add(ds, 7)获取自动调度的目标日期,手动触发时用传入参数覆盖:
def process_s3_file(**context): dag_run_conf = context.get('dag_run').conf or {} target_date = dag_run_conf.get('target_date') if not target_date: # 自动调度时,将logical_date(上周一)偏移一周得到当前周一日期 target_date = context['macros'].ds_add(context['ds'], 7) # 文件处理逻辑
自动调度时会自动匹配当前触发日的文件,手动触发时传入指定日期即可处理对应文件。
方案3:改用S3传感器触发(适合文件延迟频繁场景)
如果文件到达时间不稳定,固定cron调度不合适,可用S3KeySensor监听文件,到达后自动触发处理:
实现步骤:
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from airflow.decorators import task from airflow.models import DAG from datetime import datetime @task def process_s3_file(**context): dag_run_conf = context.get('dag_run').conf or {} target_date = dag_run_conf.get('target_date') if target_date: file_key = f"path/to/file_{target_date}.csv" else: # 自动触发时取传感器检测到的文件路径 file_key = context['task_instance'].xcom_pull(task_ids='wait_for_s3_file') target_date = file_key.split('_')[-1].split('.')[0] print(f"Processing file: {file_key}") # 文件处理逻辑 with DAG( dag_id='s3_file_processor_sensor', start_date=datetime(2023, 11, 6), schedule_interval=None, # 不使用固定调度,由传感器触发 catchup=False, ) as dag: wait_for_file = S3KeySensor( task_id='wait_for_s3_file', bucket_name='your-s3-bucket', bucket_key='path/to/file_{{ ds_nodash }}.csv', # 假设文件名含日期 aws_conn_id='aws_default', poke_interval=300, # 每5分钟检测一次 ) wait_for_file >> process_s3_file()
这种方式会自动处理到达的文件,手动触发时传入target_date也能处理指定文件。
内容的提问来源于stack exchange,提问作者timmanna
相关产品推荐
相关产品推荐

