如何让Airflow时区感知DAG按本地时区调度并获取正确日期变量?
解决Airflow悉尼时区调度下SQL模板日期变量问题
方案1:自定义模板过滤器(推荐用于SQL模板场景)
- 编写时区转换函数,用
pendulum自动处理悉尼夏令时:
import pendulum def convert_to_sydney_date(dt): sydney_tz = pendulum.timezone('Australia/Sydney') sydney_dt = dt.astimezone(sydney_tz) return sydney_dt.strftime('%Y-%m-%d')
- 在DAG中注册模板过滤器,让SQL模板可直接调用:
from airflow import DAG with DAG( dag_id='sydney_scheduled_dag', schedule_interval='0 22 * * *', # UTC 22点对应悉尼夏令时次日9点 start_date=pendulum.datetime(2023, 11, 28, tz='Australia/Sydney'), catchup=False, default_args={'owner': 'airflow'}, ) as dag: dag.template_env.filters['sydney_date'] = convert_to_sydney_date
- SQL模板中直接使用过滤器转换日期变量:
SELECT * FROM data_partition_table WHERE partition_date = '{{ data_interval_end | sydney_date }}'
方案2:通过Python任务传递转换后的日期参数
- 在Python任务中计算悉尼时区日期,通过XCom传递给SQL模板:
from airflow.operators.python import PythonOperator from airflow.operators.sql import SQLExecuteQueryOperator def calculate_sydney_dates(**context): execution_dt = context['execution_date'] sydney_tz = pendulum.timezone('Australia/Sydney') sydney_ds = execution_dt.astimezone(sydney_tz).strftime('%Y-%m-%d') sydney_next_ds = execution_dt.add(days=1).astimezone(sydney_tz).strftime('%Y-%m-%d') return {'sydney_ds': sydney_ds, 'sydney_next_ds': sydney_next_ds} get_dates_task = PythonOperator( task_id='get_sydney_dates', python_callable=calculate_sydney_dates, provide_context=True, do_xcom_push=True, ) sql_query_task = SQLExecuteQueryOperator( task_id='run_partition_query', sql=''' SELECT * FROM data_partition_table WHERE partition_date = '{{ ti.xcom_pull(task_ids='get_sydney_dates')['sydney_ds'] }}' ''', conn_id='your_db_connection', ) get_dates_task >> sql_query_task
方案3:在SQL中直接转换时区(依赖数据库支持)
利用数据库自带的时区转换函数,将UTC日期转为悉尼时区日期:
- BigQuery示例(自动处理夏令时):
SELECT * FROM data_partition_table WHERE partition_date = DATE(DATETIME('{{ ds }}', 'UTC'), 'Australia/Sydney')
- PostgreSQL示例:
SELECT * FROM data_partition_table WHERE partition_date = (('{{ ds }}'::timestamp AT TIME ZONE 'UTC') AT TIME ZONE 'Australia/Sydney')::date
关键注意事项
- Airflow核心调度基于UTC,
schedule_interval需对应悉尼9点的UTC时间:夏令时设为0 22 * * *,非夏令时设为0 23 * * *,也可通过pendulum自动生成:
schedule_interval=pendulum.daily.at('09:00').in_timezone('Australia/Sydney').as_utc().strftime('%M %H * * *')
- 所有日期处理尽量使用带时区感知的
pendulum对象,避免原生datetime的时区偏差问题。
内容的提问来源于stack exchange,提问作者FarahN
相关产品推荐
相关产品推荐

