Airflow DAG中使用Connection()连接MsSQL后端报错求助
Airflow MSSQL连接问题:实例Connection后调用get_hook报错的解决办法
问题原因
你实例化的Connection对象只是内存中的临时实例,并没有被注册到Airflow的连接存储(比如元数据库)中。get_hook()方法默认会通过指定的conn_id去查找已注册的连接,找不到这个临时对象就会抛出AirflowNotFoundException。
正确做法
方法1:直接用Connection对象创建Hook(无需注册conn_id)
不需要将连接存入Airflow的连接管理系统,直接基于已实例化的Connection对象构建Hook:
from airflow.providers.microsoft.mssql.hooks.mssql import MsSqlHook from airflow.models.connection import Connection import json # 实例化Connection对象 c = Connection( conn_id="test_mssql_1", conn_type="mssql", description="connection description", host="**************", login="**************", password="*****************", schema="***********" # extra=json.dumps(dict(this_param="some val", that_param="other val*")), ) # 直接通过Connection对象创建Hook,跳过conn_id的全局查找 mssql_hook = MsSqlHook(conn=c) # 执行SQL查询 data_records = mssql_hook.get_records("Select top 10 * from SomeTable")
方法2:先注册Connection到Airflow,再用conn_id获取Hook
如果需要后续复用这个conn_id,可以先把Connection对象存入Airflow的元数据库:
from airflow.models.connection import Connection from airflow.utils.session import create_session import json # 实例化Connection对象 c = Connection( conn_id="test_mssql_1", conn_type="mssql", description="connection description", host="**************", login="**************", password="*****************", schema="***********" # extra=json.dumps(dict(this_param="some val", that_param="other val*")), ) # 通过会话将连接注册到Airflow with create_session() as session: # 先检查是否存在同名conn_id,避免重复创建 existing_conn = session.query(Connection).filter(Connection.conn_id == c.conn_id).first() if not existing_conn: session.add(c) session.commit() # 现在可以正常通过conn_id获取Hook mssql_hook = c.get_hook() data_records = mssql_hook.get_records("Select top 10 * from SomeTable")
补充说明
如果你的密钥存储在AWS Secrets Manager,也可以配置Airflow的connections_backend来自动从Secrets Manager拉取连接信息,无需硬编码凭证。但如果只是临时使用,上面两种方法更直接高效。
内容的提问来源于stack exchange,提问作者tkansara
相关产品推荐
相关产品推荐

