如何在Airflow中执行Snowflake查询输出的外键创建ALTER语句
执行Snowflake生成的外键ALTER语句解决方案
你现在用的SnowflakeQueryOperator只执行了生成ALTER语句的查询,没有处理返回的结果,所以外键创建语句不会实际生效,可以按下面的方法修改:
方案一:修改现有Operator实现一步执行
先给SnowflakeQueryOperator新增开关参数,支持执行查询返回的DDL语句,修改后的代码如下:
class SnowflakeQueryOperator(BaseOperator): template_fields = ['sql', 'params'] template_ext = ['.sql'] @apply_defaults def __init__(self, sql, params=None, warehouse=Variable.get('default_snowflake_warehouse'), execute_result: bool = False, # 新增参数,控制是否执行查询返回的语句 *args, **kwargs): super().__init__(*args, **kwargs) self.sql = sql self.params = params self.warehouse = warehouse self.execute_result = execute_result def execute(self, context): sf_hook = SnowflakeHook(warehouse=self.warehouse) self.log.info(f'Running generate query: {self.sql}') # 执行生成ALTER语句的查询,拿到游标获取结果 cur = sf_hook.execute_query(self.sql, parameters=self.params, return_cur=True) alter_statements = [row[0] for row in cur.fetchall()] self.log.info(f'共获取到{len(alter_statements)}条外键创建语句') if self.execute_result: for idx, stmt in enumerate(alter_statements): self.log.info(f'执行第{idx+1}条DDL: {stmt}') sf_hook.execute_query(stmt)
修改后你在定义任务的时候,只需要加上execute_result=True参数即可:
snp_create_foreign_keys = SnowflakeQueryOperator( task_id='create_foreign_keys', sql='queries/foreign_keys.sql', params={ 'schema': 'qtr' }, execute_result=True, retries=0)
注意事项
- 请确保你的
foreign_keys.sql查询结果只有单列的完整ALTER语句,没有多余输出字段,否则会导致读取结果报错 - Snowflake的DDL执行默认自动提交,不需要额外加提交操作
- 可以按需给DDL执行部分加异常捕获逻辑,方便出错时定位问题语句
方案二:拆分多任务(适合需要审计的场景)
如果你需要单独留存生成的ALTER语句用于审计,也可以拆成两个任务实现:
- 第一个任务执行生成ALTER的查询,将结果写入XCom
- 第二个任务读取XCom中的语句列表,逐个执行
该方案仅适合返回的ALTER语句总大小不超过Airflow XCom单条大小限制(默认48KB)的场景,数据量较大时不推荐使用。
内容的提问来源于stack exchange,提问作者KristiLuna
相关产品推荐
相关产品推荐

