Airflow中SqlSensor无法识别Postgres连接类型的问题求助
问题分析:Airflow SqlSensor 报 Unknown hook type "postgres" 错误
问题场景
在Ubuntu系统上运行Airflow 2.8.1,已安装apache-airflow-providers-postgres==5.10.0和apache-airflow-providers-common-sql==1.10.1插件,尝试用SqlSensor检查Postgres表public.intradaytrades近1天的新增行,任务配置如下:
sql_sensor = SqlSensor( task_id="sql_sensor", conn_id="pg_conn", success=_success_criteria, sql="SELECT COUNT(*) FROM public.intradaytrades WHERE timestamp > CURRENT_DATE - INTERVAL '1 day';", mode="reschedule", fail_on_empty=True, timeout=60 * 60, # 1 hour timeout for the sensor poke_interval=60 * 5, # 5 minutes between pokes )
其中pg_conn是已验证可用的Postgres类型连接,但任务执行时报错。
错误日志
[2024-02-04, 01:10:25 UTC] {base.py:83} INFO - Using connection ID 'pg_conn' for task execution. [2024-02-04, 01:10:25 UTC] {taskinstance.py:2698} ERROR - Task failed with exception Traceback (most recent call last): File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/models/taskinstance.py", line 433, in _execute_task result = execute_callable(context=context, **execute_callable_kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/sensors/base.py", line 265, in execute raise e File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/sensors/base.py", line 247, in execute poke_return = self.poke(context) ^^^^^^^^^^^^^^^^^^ File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/providers/common/sql/sensors/sql.py", line 93, in poke hook = self._get_hook() ^^^^^^^^^^^^^^^^ File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/providers/common/sql/sensors/sql.py", line 84, in _get_hook hook = conn.get_hook(hook_params=self.hook_params) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/var/lib/airflow/python_3_11_airflow/lib/python3.11/site-packages/airflow/models/connection.py", line 363, in get_hook raise AirflowException(f'Unknown hook type "{self.conn_type}"') airflow.exceptions.AirflowException: Unknown hook type "postgres"
异常现象
直接使用PostgresHook通过pg_conn执行自定义查询时,可正常获取数据。
当前依赖版本
apache-airflow==2.8.1 apache-airflow-providers-celery==3.5.2 apache-airflow-providers-common-io==1.2.0 apache-airflow-providers-common-sql==1.10.1 apache-airflow-providers-ftp==3.7.0 apache-airflow-providers-github==2.5.1 apache-airflow-providers-http==4.8.0 apache-airflow-providers-imap==3.5.0 apache-airflow-providers-microsoft-winrm==3.4.0 apache-airflow-providers-postgres==5.10.0 apache-airflow-providers-redis==3.6.0 apache-airflow-providers-slack==8.6.0 apache-airflow-providers-sqlite==3.7.0 apache-airflow-providers-ssh==3.10.0
错误原因及解决方案
核心原因
错误本质是SqlSensor(来自common-sql插件)尝试通过连接对象自动匹配Hook时,无法识别postgres类型的连接映射。尽管直接调用PostgresHook正常,但存在以下可能触发问题的场景:
- Airflow连接元数据缓存未更新:连接创建后才安装Postgres插件,导致系统未同步连接类型与Hook的关联关系。
common-sql与Postgres插件的版本兼容细节:默认Hook自动匹配逻辑在部分版本组合下无法正确关联postgres类型连接。
解决方案
显式指定Hook类
修改SqlSensor配置,直接指定PostgresHook的路径,跳过自动匹配逻辑:sql_sensor = SqlSensor( task_id="sql_sensor", conn_id="pg_conn", success=_success_criteria, sql="SELECT COUNT(*) FROM public.intradaytrades WHERE timestamp > CURRENT_DATE - INTERVAL '1 day';", mode="reschedule", fail_on_empty=True, timeout=60 * 60, poke_interval=60 * 5, hook_class='airflow.providers.postgres.hooks.postgres.PostgresHook' )重启Airflow服务
刷新连接元数据缓存,重启webserver和scheduler:sudo systemctl restart airflow-webserver sudo systemctl restart airflow-scheduler验证连接类型
在Airflow UI的连接管理页面,确认pg_conn的连接类型为Postgres(对应底层conn_type字段为postgres),若有误则修正后重新运行任务。
内容的提问来源于stack exchange,提问作者cmdel
相关产品推荐
相关产品推荐

