Airflow 2.0本地运行时SnowflakeQueryOperator任务重复执行如何解决
问题根因
你在定义SnowflakeQueryOperator的sql参数时,直接调用了SnowHook().run_fk_alter_statements(schema,query)方法,而该方法内部包含self.execute_query的实际查询执行逻辑。Airflow调度器会定期(默认30s,你本地可能配置为5s)扫描并解析所有DAG文件,每次解析都会执行DAG定义文件里的顶层代码,所以还没等到任务实际触发运行,每次DAG解析就已经在执行Snowflake查询了,才会频繁触发DUO推送。
解决方案
- 第一步:把实际查询执行逻辑挪到Operator执行阶段,不要在DAG定义时就执行。可以用PythonOperator封装整个逻辑,避免DAG解析时触发查询。
- 第二步:拆分
run_fk_alter_statements方法的「SQL生成」和「SQL执行」逻辑,DAG定义阶段只做无副作用的SQL拼接,执行阶段再调用Hook跑实际查询。
修改后的参考代码如下:
# 拆分原方法,仅做SQL生成,不执行实际查询 def generate_fk_alter_statements(self, schema, additional_fk): fk_query_path = "/fkeys.sql" with open(f'{fk_query_path}', 'r') as fd: query = fd.read() additions = [] for fk in additional_fk: additions.append(f""" or (t2.table_name = '{fk['table_name']}' and t2.column_name = '{fk['column_name']}' and t1.table_name = '{fk['ref_table_name']}' and t1.column_name = '{fk['ref_column_name']}')\n""".upper()) # 仅返回元数据查询SQL,不实际执行 return query.format(schema=schema, fks=''.join(additions))
调整DAG中的任务定义:
from airflow.operators.python import PythonOperator def run_foreign_key_logic(**context): snow_hook = SnowHook() # 生成元数据查询SQL meta_sql = snow_hook.generate_fk_alter_statements(schema, query) # 执行元数据查询获取外键修改语句 raw_out = snow_hook.execute_query(meta_sql, fetch_all=True) query_jobs = [raw_query[0] for raw_query in raw_out] # 执行所有外键修改SQL for fk_sql in query_jobs: snow_hook.execute_query(fk_sql) create_foreign_keys = PythonOperator( dag=dag, task_id='check_and_run_foreign_key_query', python_callable=run_foreign_key_logic, trigger_rule=TriggerRule.ALL_DONE )
内容的提问来源于stack exchange,提问作者KristiLuna
相关产品推荐
相关产品推荐

