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

Airflow同一DAG按不同地区当地时间运行带参数的Notebook代码

Airflow多地区定时任务实现方案

首先明确:Airflow不支持单个DAG设置多个schedule_interval,但可以通过以下两种方案满足你的需求,优先推荐第一种,更符合Airflow最佳实践。

方案一:循环生成多地区独立DAG(推荐)

通过模板化的方式,基于地区配置列表循环创建多个DAG实例,每个DAG对应一个地区,独立设置时区、定时规则和专属参数。这种方式逻辑清晰,每个地区的任务独立,方便监控和排查问题。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
import pendulum

# 定义各地区的配置信息:地区标识、时区、定时规则、专属参数
region_configs = [
    {
        "region": "us_pst",
        "timezone": "America/Los_Angeles",
        "schedule": "0 5 * * *",  # PST 凌晨5点
        "params": {"country": "US", "data_source": "pst_db"}
    },
    {
        "region": "uk_utc",
        "timezone": "UTC",
        "schedule": "0 5 * * *",  # UTC 凌晨5点
        "params": {"country": "UK", "data_source": "utc_db"}
    },
    {
        "region": "jp_jst",
        "timezone": "Asia/Tokyo",
        "schedule": "0 5 * * *",  # JST 凌晨5点
        "params": {"country": "Japan", "data_source": "jst_db"}
    }
]

# 循环生成每个地区对应的DAG
for config in region_configs:
    with DAG(
        dag_id=f"notebook_runner_{config['region']}",
        start_date=datetime(2024, 1, 1, tzinfo=pendulum.timezone(config['timezone'])),
        schedule_interval=config['schedule'],
        catchup=False,
        default_args={"owner": "airflow"},
        params=config['params']  # 注入地区专属参数
    ) as dag:
        def run_notebook(**context):
            # 从上下文获取地区参数
            params = context["params"]
            country = params["country"]
            data_source = params["data_source"]
            # 替换为你的Notebook执行逻辑(比如调用papermill执行ipynb文件)
            print(f"Executing notebook for {country} with data source: {data_source}")

        execute_notebook = PythonOperator(
            task_id=f"run_notebook_{config['region']}",
            python_callable=run_notebook,
            provide_context=True
        )

    # 将生成的DAG注册到全局命名空间,确保Airflow能识别
    globals()[f"notebook_runner_{config['region']}"] = dag

方案优势

  • 每个地区的任务独立成DAG,日志、运行状态分开,便于监控和问题排查
  • 配置集中管理,新增地区只需在region_configs里添加一条配置即可
  • 严格遵循各地区的时区定时规则,无时间窗口误差

方案二:单DAG内分支处理(不推荐,适合简单场景)

如果一定要在同一个DAG里处理所有地区任务,可以通过定时触发DAG,再通过分支任务判断当前时间是否符合某地区的运行时间,触发对应带参数的Notebook任务。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.operators.dummy import DummyOperator
from datetime import datetime
import pendulum

def check_eligible_regions(**context):
    current_utc = datetime.utcnow()
    # 转换为各时区的当前时间
    pst_time = current_utc.astimezone(pendulum.timezone("America/Los_Angeles"))
    utc_time = current_utc.astimezone(pendulum.timezone("UTC"))
    jst_time = current_utc.astimezone(pendulum.timezone("Asia/Tokyo"))

    eligible_tasks = []
    # 判断是否到达对应地区的凌晨5点(留10分钟误差窗口)
    if pst_time.hour == 5 and pst_time.minute < 10:
        eligible_tasks.append("run_notebook_us_pst")
    if utc_time.hour == 5 and utc_time.minute < 10:
        eligible_tasks.append("run_notebook_uk_utc")
    if jst_time.hour == 5 and jst_time.minute < 10:
        eligible_tasks.append("run_notebook_jp_jst")

    return eligible_tasks if eligible_tasks else "no_task_to_run"

def run_notebook(country, data_source):
    # 替换为你的Notebook执行逻辑
    print(f"Executing notebook for {country} with data source: {data_source}")

with DAG(
    dag_id="notebook_runner_multi_region",
    start_date=datetime(2024, 1, 1, tzinfo=pendulum.timezone("UTC")),
    schedule_interval="*/10 * * * *",  # 每10分钟检查一次时间
    catchup=False,
    default_args={"owner": "airflow"}
) as dag:
    check_time = BranchPythonOperator(
        task_id="check_eligible_regions",
        python_callable=check_eligible_regions,
        provide_context=True
    )

    # 各地区的Notebook执行任务
    run_us_pst = PythonOperator(
        task_id="run_notebook_us_pst",
        python_callable=run_notebook,
        op_kwargs={"country": "US", "data_source": "pst_db"}
    )

    run_uk_utc = PythonOperator(
        task_id="run_notebook_uk_utc",
        python_callable=run_notebook,
        op_kwargs={"country": "UK", "data_source": "utc_db"}
    )

    run_jp_jst = PythonOperator(
        task_id="run_notebook_jp_jst",
        python_callable=run_notebook,
        op_kwargs={"country": "Japan", "data_source": "jst_db"}
    )

    no_task = DummyOperator(task_id="no_task_to_run")

    # 设置任务依赖
    check_time >> [run_us_pst, run_uk_utc, run_jp_jst, no_task]

方案劣势

  • 需要频繁调度DAG(比如每10分钟),增加Airflow调度压力
  • 时间判断逻辑容易出现误差,需设置合理的时间窗口
  • 所有地区任务的日志混在同一个DAG里,排查问题难度大

总结

优先选择方案一,通过循环生成多个独立DAG来实现多地区定时任务,这是Airflow社区推荐的最佳实践,逻辑清晰且易于维护。新增地区时只需补充配置,无需修改核心代码。

内容的提问来源于stack exchange,提问作者priya khokher

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 17:20:26