Apache Airflow SqlSensor连接Oracle数据库报错求助
问题解决:Airflow SqlSensor连接Oracle报错“The connection type is not supported...”
报错原因
在Airflow 2.3.3版本中,SqlSensor要求关联的数据库Hook必须是DbApiHook的子类,但你使用的apache-airflow-providers-oracle 3.2.0中的OracleHook并未继承DbApiHook,导致Sensor的兼容性校验失败。虽然测试连接可以成功,但SqlSensor的执行逻辑不认可该Hook类型。
解决方案
方案一:升级版本(推荐)
升级Airflow到2.4.0及以上,同时将Oracle Provider升级到4.0.0及以上版本。这些版本中OracleHook已调整为继承DbApiHook,可以直接被SqlSensor兼容。执行以下命令升级:
pip install --upgrade apache-airflow>=2.4.0 apache-airflow-providers-oracle>=4.0.0
方案二:自定义Oracle专用Sensor
如果无法升级环境,可以自定义一个基于OracleHook的Sensor,实现和SqlSensor相同的功能:
from airflow.sensors.base import BaseSensorOperator from airflow.providers.oracle.hooks.oracle import OracleHook from airflow.utils.decorators import apply_defaults class OracleSqlSensor(BaseSensorOperator): @apply_defaults def __init__(self, sql, oracle_conn_id='oracle_default', **kwargs): super().__init__(**kwargs) self.sql = sql self.oracle_conn_id = oracle_conn_id def poke(self, context): hook = OracleHook(oracle_conn_id=self.oracle_conn_id) records = hook.get_records(self.sql) # 判断查询结果是否非空,可根据需求调整逻辑 return len(records) > 0
在DAG中替换使用自定义Sensor:
from airflow import DAG from datetime import datetime # 注意替换为自定义Sensor所在的模块路径 from your_sensor_module import OracleSqlSensor with DAG( 'oracle_check_dag', start_date=datetime(2023, 1, 1), schedule_interval='@daily', catchup=False ) as dag: check_exec_date = OracleSqlSensor( task_id='check_exec_date', sql="SELECT 1 FROM your_table WHERE exec_date = '{{ ds }}'", oracle_conn_id='your_oracle_connection_id' )
内容的提问来源于stack exchange,提问作者vlcheong
相关产品推荐
相关产品推荐

