如何在Airflow DAG上下文中获取DAG执行小时
Airflow 获取DAG执行小时参数方案
{{ ds }} 是Airflow内置的仅返回YYYY-MM-DD格式逻辑执行日期的模板变量,本身不携带时间信息,要获取调度对应的执行小时直接用内置的execution_date模板变量做格式化即可,不需要额外开发自定义逻辑。
可用写法
根据你需要的格式和Airflow版本选对应写法即可:
- 通用兼容写法(所有Airflow版本通用,返回00-23的补零两位字符串格式小时,比如凌晨2点返回
02,下午4点返回16):{{ execution_date.strftime('%H') }} - Airflow 2.0+ 简写写法(返回0-23的整数型小时值):
{{ execution_date.hour }}
注意:
execution_date对应DAG的逻辑调度时间,不是任务实际启动运行的物理时间,对你配置的每小时第30分钟触发的调度规则(30 * * * *)来说,取到的小时值完全匹配调度周期对应的小时,符合你的业务需求。如果返回小时和本地时间差8小时,在DAG初始化参数中添加timezone="Asia/Shanghai"指定东八区时区即可,Airflow默认使用UTC时区。
修改后可直接运行的示例代码
with DAG(dag_id="dag_name", schedule_interval="30 * * * *", max_active_runs=1) as dag: features_hourly = KubernetesPodOperator( task_id="task-name", name="task-name", cmds=[ "python", "-m", "sql_library.scripts.sql_executor", "--template", "format", "--env-names", "'" + json.dumps(["SCHEMA"]) + "'", "--vars", "'" + json.dumps({ "EXECUTION_DATE": "{{ ds }}", "PREDICTION_HOUR": "{{ execution_date.strftime('%H') }}", }) + "'", "sql_filename.sql", ], **default_task_params, )
内容的提问来源于stack exchange,提问作者RTM
相关产品推荐
相关产品推荐

