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
相关产品推荐
相关产品推荐

