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

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()调用即可:

  1. 自定义非查询操作的结果处理器,直接返回SQL执行影响的行数
  2. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 10:51:22