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

Celery任务结合Psycopg报错:ProgrammingError 上次操作未产生结果

问题根源分析

  1. 连接池连接未正确重置:当前get_db_conn仅负责获取和放回连接,未在放回前确保连接处于干净的事务状态。若任务因异常或未完成操作导致连接处于非IDLE状态,放回池后被其他任务复用,就会触发there is already a transaction in progress警告。
  2. 事务边界管理不严谨:任务中跨commit操作复用同一个cursor,可能导致连接状态混乱;异常处理中事务回滚后的连接状态未同步,进一步加剧了连接状态污染。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 11:17:14