升级MWAA Airflow至2.7.2后common-sql与AwsGenericHook报错解决及代码修正
问题概况
MWAA Airflow从2.0.2升级到2.7.2后,部分使用SnowflakeOperator的DAG抛出common-sql相关异常,提示AwsGenericHook不支持common-sql规范。
报错详情
airflow.exceptions.AirflowException: You are trying to use
common-sqlwith AwsGenericHook, but its provider does not support it. Please upgrade the provider to a version that supportscommon-sql. The hook class should be a subclass ofairflow.providers.common.sql.hooks.sql.DbApiHook. Got AwsGenericHook Hook with class hierarchy: [<class 'airflow.providers.amazon.aws.hooks.base_aws.AwsGenericHook'>, <class 'airflow.hooks.base.BaseHook'>, <class 'airflow.utils.log.logging_mixin.LoggingMixin'>, <class 'typing.Generic'>, <class 'object'>]
关键报错堆栈:
File "/usr/local/airflow/.local/lib/python3.11/site-packages/airflow/providers/common/sql/operators/sql.py", line 167, in _hook raise AirflowException(
当前requirements.txt配置
--constraint "https://raw.githubusercontent.com/apache/airflow/constraints-2.7.2/constraints-3.11.txt" apache-airflow-providers-amazon apache-airflow-providers-ftp apache-airflow-providers-slack apache-airflow-providers-imap apache-airflow-providers-sqlite apache-airflow-providers-snowflake apache-airflow-providers-mysql snowflake-connector-python snowflake-sqlalchemy slack-sdk pandas boto3 botocore pendulum pytz pytzdata imapclient xlsxwriter pydrive openpyxl cryptography
问题代码片段
process_csv = SnowflakeOperator( task_id=f'process_csv_{table_name}', sql=f'create or replace table "{snowflake__database}"."{snowflake__schema}".{table_name} clone "{snowflake__database}"."{snowflake__schema}".STG_{table_name}', snowflake_conn_id="snowflake_conn_id", retries=0, on_failure_callback=task_fail_slack_alert )
修复方案
1. 修正Snowflake连接配置(核心解决)
报错显示Airflow错误调用了AwsGenericHook,说明snowflake_conn_id对应的连接类型配置错误:
- 登录Airflow UI,进入
Admin > Connections - 找到
snowflake_conn_id连接,确认Conn Type为Snowflake(而非AWS类连接) - 核对连接的账户、数据库、仓库、角色等参数是否正确
2. 锁定兼容的Snowflake依赖版本
Airflow 2.7.2要求snowflake provider需支持common-sql规范,在requirements.txt中明确指定匹配约束文件的版本:
apache-airflow-providers-snowflake==5.1.0 snowflake-connector-python==3.3.0 snowflake-sqlalchemy==1.5.0
(版本号需与Airflow 2.7.2的constraints文件保持一致,避免版本冲突)
3. 代码显式指定Hook类(兜底方案)
若连接配置无误仍报错,可在SnowflakeOperator中强制指定Snowflake Hook类:
from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook process_csv = SnowflakeOperator( task_id=f'process_csv_{table_name}', sql=f'create or replace table "{snowflake__database}"."{snowflake__schema}".{table_name} clone "{snowflake__database}"."{snowflake__schema}".STG_{table_name}', snowflake_conn_id="snowflake_conn_id", hook_class=SnowflakeHook, # 显式指定Hook类型 retries=0, on_failure_callback=task_fail_slack_alert )
验证步骤
- 更新MWAA的
requirements.txt并重新部署环境 - 确认Airflow连接配置正确
- 触发目标DAG,验证报错消除
内容的提问来源于stack exchange,提问作者Xi12

