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

Apache Airflow同任务生成多实例问题求助

Apache Airflow任务生成多实例问题解决

我是Apache Airflow新手,编写的DAG为每个省份创建了一个PythonOperator任务,但DAG运行时每个任务会生成多个实例(例如省份1的任务生成5个实例),期望每个任务仅对应一个实例。相关代码如下:

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "start_date": datetime(2023, 7, 28),
    "email": ["airflow@airflow.com"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}


def run_building_pp_checks(province_id):
    query = f'''
       with building_section as(
       select id,population_point_id,province_id,county_id,zone_id,rural_id from building
       where province_id = {province_id} and deleted_at is null

       )

       select building.id from public.population_point pp
       join building_section building on building.population_point_id = pp.id and (pp.province_id != {province_id} or building.county_id != pp.county_id or  building.rural_id != pp.rural_id
        or building.zone_id != pp.zone_id ) 
       '''
    print(query)
    postgres_hook = PostgresHook(
            postgres_conn_id='my_postgres_conn',
         )
    result=postgres_hook.get_records(query)
    print(result)


    if result :
        values=[(id,'building',datetime.now().strftime("%Y-%m-%d %H:%M:%S")) for id in result]
        print(values)
        postgres_hook.insert_rows(table='report.mismatch_records', rows=values , target_fields=['id' , 'table_name' , 'inserted_at'])


dag1 = DAG(
    'building_pp_checks',
    default_args=default_args,
    # 此处原代码未显式定义schedule_interval,默认可能为@daily或其他值
)

for province_id in range(1,39):
        task=PythonOperator(
            task_id=f'task_province_{province_id}',
            python_callable=run_building_pp_checks,
            dag=dag1,
            op_kwargs={'province_id': province_id},
        )

        if province_id!=1:
            task.set_upstream(previous_task)
        previous_task=task

注:start_date设置为运行时间的两天前。


问题原因

核心原因是DAG调度配置与start_date的组合触发Airflow自动补全历史运行:

  • 若DAG设置了schedule_interval(非None,默认可能为@daily),且start_date早于当前时间,Airflow会自动生成从start_date到当前时间的所有调度间隔对应的DAG Run,每个DAG Run都会触发对应任务的实例,因此出现多实例。
  • 即使depends_on_past=False,也无法阻止Airflow生成这些历史补全的运行实例。

解决方案

根据你的需求选择以下方案:

1. 一次性运行的DAG(无需定时调度)

将DAG的schedule_interval设置为None,Airflow不会自动生成任何调度运行,仅支持手动触发一次:

dag1 = DAG(
    'building_pp_checks',
    default_args=default_args,
    schedule_interval=None,
)

2. 需要定时调度,但跳过历史补全运行

在DAG定义中添加catchup=False参数,Airflow会跳过start_date到当前时间的所有历史调度间隔,仅从下一个调度时间开始运行:

dag1 = DAG(
    'building_pp_checks',
    default_args=default_args,
    schedule_interval='@daily',  # 替换为你的实际调度间隔,如@hourly等
    catchup=False,
)

3. 清理已生成的多余实例

若已经生成了多余的任务实例,可在Airflow UI中操作:

  • 进入对应DAG页面,切换到「Tree View」或「Graph View」
  • 选中多余的DAG Run,执行「Clear」操作,删除不必要的运行实例

额外优化建议

原代码中SQL查询使用字符串拼接存在SQL注入风险,建议改用参数化查询:

def run_building_pp_checks(province_id):
    query = '''
       with building_section as(
       select id,population_point_id,province_id,county_id,zone_id,rural_id from building
       where province_id = %s and deleted_at is null

       )

       select building.id from public.population_point pp
       join building_section building on building.population_point_id = pp.id and (pp.province_id != %s or building.county_id != pp.county_id or  building.rural_id != pp.rural_id
        or building.zone_id != pp.zone_id ) 
       '''
    print(query)
    postgres_hook = PostgresHook(
            postgres_conn_id='my_postgres_conn',
         )
    # 使用parameters传递参数,避免SQL注入
    result=postgres_hook.get_records(query, parameters=(province_id, province_id))
    print(result)

    # 后续逻辑保持不变...

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 11:05:36