Airflow操作DB2的DAG报jaydebeapi.Error但数据执行成功问题求解
Airflow JDBC连接DB2执行写操作抛jaydebeapi.Error异常问题
问题现象
参照Airflow官方文档编写连接DB2的Airflow DAG,执行数据插入/更新操作时会抛出jaydebeapi.Error异常:DB2侧数据实际已完成插入/更新,但Airflow UI中对应DAG任务会被标记为FAILED。升级JayDeBeApi、JPype1依赖,降级Airflow版本均无法解决问题。
问题复现代码
with DAG("my_dag1", default_args=default_args, schedule_interval="@daily", catchup=False) as dag: cerating_table = JdbcOperator( task_id='creating_table', jdbc_conn_id='db2', sql=r""" insert into DB2ECIF.T2(C1,C1_DATE) VALUES('TEST',CURRENT DATE); """, autocommit=True, dag=dag )
报错日志
[2022-06-20 02:16:03,743] {base.py:68} INFO - Using connection ID 'db2' for task execution. [2022-06-20 02:16:04,785] {dbapi.py:213} INFO - Running statement: insert into DB2ECIF.T2(C1,C1_DATE) VALUES('TEST',CURRENT DATE); , parameters: None [2022-06-20 02:16:04,842] {dbapi.py:221} INFO - Rows affected: 1 [2022-06-20 02:16:04,844] {taskinstance.py:1889} ERROR - Task failed with exception Traceback (most recent call last): File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/jdbc/operators/jdbc.py", line 76, in execute return hook.run(self.sql, self.autocommit, parameters=self.parameters, handler=fetch_all_handler) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/hooks/dbapi.py", line 195, in run result = handler(cur) File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/jdbc/operators/jdbc.py", line 30, in fetch_all_handler return cursor.fetchall() File "/home/airflow/.local/lib/python3.7/site-packages/jaydebeapi/__init__.py", line 596, in fetchall row = self.fetchone() File "/home/airflow/.local/lib/python3.7/site-packages/jaydebeapi/__init__.py", line 561, in fetchone raise Error() jaydebeapi.Error [2022-06-20 02:16:04,847] {taskinstance.py:1400} INFO - Marking task as FAILED. dag_id=my_dag1, task_id=creating_table, execution_date=20210101T000000, start_date=, end_date=2022-06-20T02:16:04
环境版本信息
- Airflow/2.3.2
- IBM DB2/11.5.7
- OpenJDK/15.0.2
- JayDeBeApi/1.2.0(已尝试升级至1.2.3,问题未解决)
- JPype1/0.7.2(已尝试升级至1.4.0,问题未解决)
- apache-airflow-providers-jdbc/3.0.0
- 已尝试降级Airflow至2.2.3、2.2.5,问题仍复现
根因分析
从报错栈可以定位到异常触发点:JdbcOperator默认内置了fetch_all_handler结果处理器,SQL执行完成后会自动调用cursor.fetchall()拉取全量结果集,这个逻辑是为SELECT类查询场景设计的。
INSERT/UPDATE/DELETE/DDL这类写操作执行完成后不会返回可查询的结果集,JayDeBeApi驱动对无结果集的游标调用fetchall()方法时,会直接抛出无附加信息的jaydebeapi.Error。日志中Rows affected: 1的打印已经证明SQL本身在DB2侧执行成功,异常是执行完SQL后的结果集拉取步骤触发的,和依赖版本、Airflow版本无关。
解决方案
针对写操作替换默认的结果处理器,跳过无意义的fetchall()调用即可:
- 自定义非查询操作的结果处理器,直接返回SQL执行影响的行数
- 在
JdbcOperator初始化时传入自定义处理器,覆盖默认配置
修复后的代码示例:
# 自定义非查询操作处理器,不执行fetchall,直接返回受影响行数 def non_query_handler(cursor): return cursor.rowcount with DAG("my_dag1", default_args=default_args, schedule_interval="@daily", catchup=False) as dag: creating_table = JdbcOperator( task_id='creating_table', jdbc_conn_id='db2', sql=r""" insert into DB2ECIF.T2(C1,C1_DATE) VALUES('TEST',CURRENT DATE); """, autocommit=True, handler=non_query_handler, # 替换默认的fetch_all_handler dag=dag )
注意:SELECT类查询场景仍可使用默认配置,仅写操作需要替换处理器。
内容的提问来源于stack exchange,提问作者Allen
相关产品推荐
相关产品推荐

