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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 12:32:09