如何在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
相关产品推荐
相关产品推荐

