Airflow2.1.2使用PostgresOperator连Redshift报递归深度超限如何解决
问题根因
你的报错核心原因有两个:
- Python环境存在
psycopg2相关包的冲突:日志里同时出现psycopg2和psycopg2abc两个包的connect方法,二者互相调用触发了递归溢出。 - 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
相关产品推荐
相关产品推荐

