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

使用SQLAlchemy/psycopg2时For循环无法更新Pandas DataFrame的问题

问题:循环查询PostgreSQL无法获取最新数据的原因与解决办法

我编写了一个For循环,期望每隔5秒从被其他线程每5秒更新一次的PostgreSQL表中获取数据,以更新Pandas DataFrame。测试发现,不使用循环时能获取到最新的更新时间,但加入循环后结果始终停留在首次查询的值。同时尝试psycopg2库也存在相同问题,请问该问题的原因是什么,如何解决?

SQLAlchemy实现代码

metadata = MetaData(bind=None)
table = Table(
    'datastore', 
    metadata, 
    autoload=True, 
    autoload_with=engine
)

stmt = select([
    table.columns.date,
    table.columns.open,
    table.columns.high,
    table.columns.low,
    table.columns.close
]).where(and_(table.columns.date_ == datetime.today().strftime('%Y-%m-%d') and table.columns.close != 0))
#]).where(and_(table.columns.date_ == '2023-01-12' and table.columns.close != 0))
    
connection = engine.connect()

for x in range(1000000):
    data_from_db = pd.DataFrame(connection.execute(stmt).fetchall())  
    data_from_db = data_from_db[data_from_db['close'] != 0]
    print(data_from_db.date.iloc[-1])
    time.sleep(5)

psycopg2实现代码

for x in range(1000000):
    conn = psycopg2.connect(
                           host='localhost', 
                           database='ABC',  
                           user='postgres',
                           password='*******')
    cur = conn.cursor()
    cur.execute("select max(date) from public.datastore")

    y = cur.fetchall()
    print(y)
    time.sleep(5)

问题原因

  1. SQLAlchemy 场景:

    • 你在循环外创建了单个连接,所有循环内的查询都在同一个事务中执行。PostgreSQL默认事务隔离级别为READ COMMITTED,事务一旦开启,会基于启动时的数据快照执行查询,无法看到其他事务后续提交的新数据。
    • 另外你的SQL语句中,datetime.today().strftime('%Y-%m-%d')是在循环外计算的,只会取一次当天日期(不过这不是当前问题的核心)。
  2. psycopg2 场景:

    • psycopg2 默认关闭自动提交(autocommit=False),每次创建连接后执行查询会隐式开启事务,若未显式提交/关闭事务,新连接的查询依然无法获取最新提交的数据;同时你每次循环创建新连接却未关闭,可能导致连接泄漏,进一步影响数据可见性。

解决办法

针对 SQLAlchemy 的修复

有两种可行方案:

  • 方案1:每次循环创建新连接
    将连接创建逻辑放入循环内部,确保每次查询都在独立事务中执行,同时用with语句自动管理连接生命周期:

    metadata = MetaData(bind=None)
    table = Table(
        'datastore', 
        metadata, 
        autoload=True, 
        autoload_with=engine
    )
    
    for x in range(1000000):
        # 每次循环重新生成日期条件,避免日期固定
        today_date = datetime.today().strftime('%Y-%m-%d')
        stmt = select([
            table.columns.date,
            table.columns.open,
            table.columns.high,
            table.columns.low,
            table.columns.close
        ]).where(and_(table.columns.date_ == today_date, table.columns.close != 0))
        
        # 每次循环创建新连接,with语句自动关闭连接
        with engine.connect() as connection:
            data_from_db = pd.DataFrame(connection.execute(stmt).fetchall())  
            data_from_db = data_from_db[data_from_db['close'] != 0]
            print(data_from_db.date.iloc[-1])
        time.sleep(5)
    
  • 方案2:复用连接并手动提交事务
    如果要保留单个连接,每次查询后提交事务,强制更新数据快照:

    metadata = MetaData(bind=None)
    table = Table(
        'datastore', 
        metadata, 
        autoload=True, 
        autoload_with=engine
    )
    
    today_date = datetime.today().strftime('%Y-%m-%d')
    stmt = select([
        table.columns.date,
        table.columns.open,
        table.columns.high,
        table.columns.low,
        table.columns.close
    ]).where(and_(table.columns.date_ == today_date, table.columns.close != 0))
        
    connection = engine.connect()
    try:
        for x in range(1000000):
            data_from_db = pd.DataFrame(connection.execute(stmt).fetchall())  
            data_from_db = data_from_db[data_from_db['close'] != 0]
            print(data_from_db.date.iloc[-1])
            # 提交事务,更新数据快照
            connection.commit()
            time.sleep(5)
    finally:
        # 确保连接最终关闭
        connection.close()
    

针对 psycopg2 的修复

开启自动提交,或者每次查询后手动提交事务,并严格关闭游标与连接:

for x in range(1000000):
    conn = psycopg2.connect(
        host='localhost', 
        database='ABC',  
        user='postgres',
        password='*******'
    )
    # 开启自动提交,确保查询能获取最新数据
    conn.autocommit = True
    cur = conn.cursor()
    cur.execute("select max(date) from public.datastore")

    y = cur.fetchall()
    print(y)
    # 关闭游标和连接,避免资源泄漏
    cur.close()
    conn.close()
    time.sleep(5)

也可以选择手动提交事务的方式:

for x in range(1000000):
    conn = psycopg2.connect(
        host='localhost', 
        database='ABC',  
        user='postgres',
        password='*******'
    )
    cur = conn.cursor()
    cur.execute("select max(date) from public.datastore")

    y = cur.fetchall()
    print(y)
    # 手动提交事务
    conn.commit()
    cur.close()
    conn.close()
    time.sleep(5)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 05:20:49