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

Airflow周/月调度:以处理日期为结束日期时如何选起始日期

解决Airflow周调度DAG处理对应周期文件的问题

方案1:结合data_interval_end宏与手动触发参数

这是最适配你需求的方案:自动调度时用data_interval_end匹配文件到达日期,文件延迟时通过手动传入参数指定目标日期。

实现步骤:

  1. 任务函数优先读取手动触发传入的配置参数,无参数时使用周期结束日期:
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
    )
  1. 文件延迟时手动触发:
    在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 12:57:06