You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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正常,但存在以下可能触发问题的场景:

  1. Airflow连接元数据缓存未更新:连接创建后才安装Postgres插件,导致系统未同步连接类型与Hook的关联关系。
  2. common-sql与Postgres插件的版本兼容细节:默认Hook自动匹配逻辑在部分版本组合下无法正确关联postgres类型连接。

解决方案

  1. 显式指定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'
    )
    
  2. 重启Airflow服务
    刷新连接元数据缓存,重启webserver和scheduler:

    sudo systemctl restart airflow-webserver
    sudo systemctl restart airflow-scheduler
    
  3. 验证连接类型
    在Airflow UI的连接管理页面,确认pg_conn的连接类型为Postgres(对应底层conn_type字段为postgres),若有误则修正后重新运行任务。

内容的提问来源于stack exchange,提问作者cmdel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 14:29:58