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

