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

Airflow2.1.2使用PostgresOperator连Redshift报递归深度超限如何解决

问题根因

你的报错核心原因有两个:

  1. Python环境存在psycopg2相关包的冲突:日志里同时出现psycopg2和psycopg2abc两个包的connect方法,二者互相调用触发了递归溢出。
  2. DAG代码存在缩进错误:DAG实例的定义被缩进在load_data_to_redshift函数内部,作用域异常,同时函数本身也存在多余的前置缩进,可能引发非预期的变量引用问题。
解决方案
  • 清理并重新安装psycopg2依赖:
    执行以下命令卸载所有冲突的psycopg2包,再安装兼容Python3.8和Airflow2.1.2的稳定版本:
pip uninstall psycopg2 psycopg2-binary psycopg2abc -y
pip install psycopg2-binary==2.9.3
  • 修正DAG代码的缩进和作用域问题:
    修正后的DAG代码参考如下:
import datetime
import logging

from airflow import DAG
from airflow.contrib.hooks.aws_hook import AwsHook
from airflow.hooks.postgres_hook import PostgresHook
from airflow.operators.postgres_operator import PostgresOperator
from airflow.operators.python_operator import PythonOperator

import sql_statements

# 函数顶格定义,去掉多余前置缩进
def load_data_to_redshift(*args, **kwargs):
    aws_hook = AwsHook("aws_credentials")
    credentials = aws_hook.get_credentials()
    # 连接id和Operator保持一致,如果你后台配置的是redshift_default就用这个
    redshift_hook = PostgresHook("redshift_default")
    sql_stmt = sql_statements.COPY_ALL_data_SQL.format(
        credentials.access_key,
        credentials.secret_key,
    )
    redshift_hook.run(sql_stmt)

# DAG定义移到函数外部,顶格写,作用域为全局
dag = DAG(
    'exercise1',
    start_date=datetime.datetime.now()
)

create_t1_table = PostgresOperator(
    task_id="create_t1_table",
    dag=dag,
    postgres_conn_id="redshift_default",
    sql=sql_statements.CREATE_t1_TABLE_SQL
)
    
create_t2_table = PostgresOperator(
    task_id="create_t2_table",
    dag=dag,
    postgres_conn_id="redshift_default",
    sql=sql_statements.CREATE_t2_TABLE_SQL,
)

create_t1_table >> create_t2_table
  • 校验Airflow连接配置:登录Airflow后台的连接管理页面,确认redshift_default的配置项完全正确,没有残留的本地Postgres连接参数。
增强日志排查的方法

如果执行上述操作后仍有问题,可以通过以下方式获取更详细的排查信息:

  • 新增测试任务打印连接配置:在DAG中添加一个临时PythonOperator,打印连接的完整信息确认是否符合预期:
def check_conn_config(**context):
    conn = PostgresHook.get_connection("redshift_default")
    logging.info(f"连接配置:host={conn.host}, port={conn.port}, schema={conn.schema}, login={conn.login}")

check_conn = PythonOperator(
    task_id="check_conn",
    dag=dag,
    python_callable=check_conn_config,
    provide_context=True
)

# 把这个任务放在建表任务前面执行
check_conn >> create_t1_table >> create_t2_table
  • 打印调用栈定位递归来源:在/home/8085/.local/lib/python3.8/site-packages/psycopg2/__init__.py的connect方法开头添加调用栈打印代码,复现问题时就能看到完整的递归调用链路:
import traceback
traceback.print_stack()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 02:15:00