Celery任务结合Psycopg报错:ProgrammingError 上次操作未产生结果
问题根源分析
- 连接池连接未正确重置:当前
get_db_conn仅负责获取和放回连接,未在放回前确保连接处于干净的事务状态。若任务因异常或未完成操作导致连接处于非IDLE状态,放回池后被其他任务复用,就会触发there is already a transaction in progress警告。 - 事务边界管理不严谨:任务中跨
commit操作复用同一个cursor,可能导致连接状态混乱;异常处理中事务回滚后的连接状态未同步,进一步加剧了连接状态污染。 - fetchone报错关联:当连接处于未完成事务中时,后续SELECT执行可能因事务上下文异常导致无结果返回,触发
the last operation didn't produce a result错误。
解决方案
1. 修复连接池的连接回收逻辑
修改get_db_conn上下文管理器,确保放回池的连接处于干净的IDLE状态:
import os import psycopg from psycopg_pool import ConnectionPool from contextlib import contextmanager PG_USERNAME = os.getenv('PG_USERNAME') if not PG_USERNAME: raise ValueError(f"Invalid postgres username") PG_PASSWORD = os.getenv('PG_PASSWORD') if not PG_PASSWORD: raise ValueError(f"Invalid postgres pass") PG_HOST = os.getenv('PG_HOST') if not PG_HOST: raise ValueError(f"Invalid postgres host") PG_PORT = os.getenv('PG_PORT') if not PG_PORT: raise ValueError(f"Invalid postgres port") conninfo = f'host={PG_HOST} port={PG_PORT} dbname=postgres user={PG_USERNAME} password={PG_PASSWORD}' connection_pool = ConnectionPool( min_size=4, max_size=100, conninfo=conninfo, check=ConnectionPool.check_connection, ) @contextmanager def get_db_conn(): conn = connection_pool.getconn() try: # 获取连接时先清理残留事务 if conn.status != psycopg.pq.ConnStatus.IDLE: conn.rollback() yield conn finally: try: # 放回前强制回滚未提交事务,重置连接会话 if conn.status != psycopg.pq.ConnStatus.IDLE: conn.rollback() conn.reset() except Exception: # 重置失败则直接关闭连接,避免污染连接池 connection_pool.close(conn) else: connection_pool.putconn(conn)
2. 优化任务中的事务与cursor管理
拆分事务边界,避免跨commit复用cursor,确保每个操作的事务状态清晰:
@app.task(bind=True) def example_task(self, id): with get_db_conn() as conn: try: # 第一个事务:查询并更新状态为running with conn.cursor(row_factory=dict_row) as cursor: cursor.execute('SELECT * FROM test WHERE id = %s', (id,)) test = cursor.fetchone() if not test: logger.warning(f'Test entry {id} not found') return cursor.execute("UPDATE test SET status = 'running' WHERE id = %s", (id,)) conn.commit() # 执行中间处理逻辑 # Some processing... # 第二个事务:获取资源并更新结果 with conn.cursor(row_factory=dict_row) as cursor: cursor.execute('SELECT * FROM test WHERE id = %s', (test['resource_id'],)) resource = cursor.fetchone() if not resource: logger.warning(f'Resource {test["resource_id"]} not found') with conn.cursor(row_factory=dict_row) as err_cursor: err_cursor.execute(""" UPDATE test SET status = 'error', error = %s WHERE id = %s """, (Jsonb({'error': 'Resource not found'}), id)) conn.commit() return cursor.execute(""" UPDATE test SET status = 'done', properties = %s WHERE id = %s """, (Jsonb(properties), id)) conn.commit() except Exception as e: logger.exception(f'Error processing task {id}: {e}') with conn.cursor(row_factory=dict_row) as cursor: cursor.execute(""" UPDATE test SET status = 'error', error = %s WHERE id = %s """, (Jsonb({'error': str(e), 'stacktrace': traceback.format_exc()}), id)) conn.commit()
3. 可选:自定义连接池检查机制
添加更严格的连接检查,确保从池取出的连接完全可用:
def custom_connection_check(conn): # 强制清理残留事务 if conn.status != psycopg.pq.ConnStatus.IDLE: conn.rollback() # 执行测试查询验证连接可用性 with conn.cursor() as cur: cur.execute("SELECT 1") return True connection_pool = ConnectionPool( min_size=4, max_size=100, conninfo=conninfo, check=custom_connection_check, )
关键说明
- 连接重置的必要性:psycopg默认autocommit=False,会自动开启事务。若任务未显式提交/回滚就放回连接,事务会持续处于打开状态,导致后续任务复用连接时出现冲突。
- 事务边界清晰化:将不同操作拆分为独立事务块,避免跨commit复用cursor,减少状态混乱的概率。
- 异常安全的回收:在finally块中处理连接重置,即使任务抛出异常,也能确保连接干净地回到池内。
内容的提问来源于stack exchange,提问作者Javierd98
相关产品推荐
相关产品推荐

