Airflow技术问题:如何在SQL查询中传入prev_execution_date减去2小时的值
在Airflow中调整prev_execution_date并传入SQL查询的方法
当然有直接实现的办法!Airflow的Jinja2模板支持对prev_execution_date这类datetime对象直接进行时间运算,不需要额外写复杂的逻辑,下面给你两种常用的实现方式:
方法1:直接在SQL模板中计算调整后的时间
因为prev_execution_date本身就是一个Python datetime对象,你可以借助Airflow内置的macros.timedelta直接在模板里减去2小时,然后转换成你需要的字符串格式。
修正后的SQL查询语句可以写成这样(优化了原语句的引号嵌套,避免混乱):
SELECT * FROM {{ mysql_table_name }} WHERE {{ incremental_column }} >= '{{ (prev_execution_date - macros.timedelta(hours=2)).to_datetime_string() }}' AND {{ incremental_column }} < '{{ execution_date.to_datetime_string() }}';
- 这里
macros.timedelta是Airflow预定义的宏,不需要额外导入,直接在模板里用就行 - 如果你的字段只需要日期部分(不需要时分秒),可以把
to_datetime_string()换成ds过滤器,比如{{ (prev_execution_date - macros.timedelta(hours=2)) | ds }}
方法2:在Operator参数中预定义调整后的值
如果觉得SQL模板里的计算逻辑太冗长,你也可以在Operator的params参数里提前定义好调整后的日期,让SQL语句更简洁:
from airflow.operators.sql import SqlOperator adjusted_query_task = SqlOperator( task_id="run_adjusted_incremental_query", sql=""" SELECT * FROM {{ params.table_name }} WHERE {{ params.inc_col }} >= '{{ params.adjusted_prev_date }}' AND {{ params.inc_col }} < '{{ execution_date.to_datetime_string() }}'; """, params={ "table_name": "your_target_table", "inc_col": "your_incremental_column", "adjusted_prev_date": "{{ (prev_execution_date - macros.timedelta(hours=2)).to_datetime_string() }}" }, conn_id="your_mysql_conn_id", dag=your_dag_object )
小提醒
如果你的数据库时区和Airflow的执行时区(默认UTC)不一致,记得在时间转换时考虑时区偏移,比如用macros.pendulum来做时区转换,避免数据范围错误。
内容的提问来源于stack exchange,提问作者pythonlearner
相关产品推荐
相关产品推荐

