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

如何在Airflow的PythonOperator中传入execution_date与prev_execution_date?

用PythonOperator替代BashOperator实现带日期模板的任务

要把你原来的BashOperator逻辑改成PythonOperator,有两种常用方式获取execution_date和next_execution_date,下面分别给出实现代码:

方法一:通过上下文参数直接获取日期

这种方式让Python函数接收Airflow的上下文参数,从中提取所需的日期变量。

首先导入必要的模块:

import os
import requests
from datetime import datetime as dt
from airflow import DAG
from airflow.operators.python import PythonOperator

定义DAG结构(和原代码一致):

dag = DAG(
    dag_id="06_templated_query",
    schedule_interval="@daily",
    start_date=dt(year=2019, month=1, day=1),
    end_date=dt(year=2019, month=1, day=5),
)

编写处理任务的Python函数,通过**kwargs接收上下文:

def fetch_events_func(**kwargs):
    # 从上下文提取执行日期和下一次执行日期
    execution_date = kwargs["execution_date"]
    next_execution_date = kwargs["next_execution_date"]
    
    # 格式化为YYYY-MM-DD字符串
    start_date_str = execution_date.strftime("%Y-%m-%d")
    end_date_str = next_execution_date.strftime("%Y-%m-%d")
    
    # 创建目标目录(exist_ok=True避免目录已存在时报错)
    os.makedirs("/data/events", exist_ok=True)
    
    # 发送HTTP请求并保存结果
    api_url = f"http://events_api:5000/events?start_date={start_date_str}&end_date={end_date_str}"
    response = requests.get(api_url)
    response.raise_for_status()  # 若请求失败抛出异常
    
    with open("/data/events.json", "w") as f:
        f.write(response.text)

创建PythonOperator实例,开启provide_context=True以传递上下文:

fetch_events = PythonOperator(
    task_id="fetch_events",
    python_callable=fetch_events_func,
    provide_context=True,
    dag=dag,
)

方法二:通过op_kwargs传入模板化日期字符串

这种方式直接在op_kwargs中使用Airflow的模板语法,把格式化后的日期字符串传入函数,函数不需要处理上下文。

同样先导入模块、定义DAG,然后编写简化的函数:

def fetch_events_func(start_date, end_date):
    os.makedirs("/data/events", exist_ok=True)
    
    api_url = f"http://events_api:5000/events?start_date={start_date}&end_date={end_date}"
    response = requests.get(api_url)
    response.raise_for_status()
    
    with open("/data/events.json", "w") as f:
        f.write(response.text)

创建PythonOperator时,在op_kwargs中使用模板语法传入日期:

fetch_events = PythonOperator(
    task_id="fetch_events",
    python_callable=fetch_events_func,
    op_kwargs={
        "start_date": "{{execution_date.strftime('%Y-%m-%d')}}",
        "end_date": "{{next_execution_date.strftime('%Y-%m-%d')}}"
    },
    dag=dag,
)

注意事项

  • 确保你的Airflow环境安装了requests库,若未安装可通过pip install requests添加。
  • Airflow 2.x中provide_context=True依然有效,也可以使用templates_dict等其他方式传递模板变量,但上述两种方法是最常用的。

内容的提问来源于stack exchange,提问作者le Minh Nguyen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 03:54:31