使用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)
问题原因
SQLAlchemy 场景:
- 你在循环外创建了单个连接,所有循环内的查询都在同一个事务中执行。PostgreSQL默认事务隔离级别为READ COMMITTED,事务一旦开启,会基于启动时的数据快照执行查询,无法看到其他事务后续提交的新数据。
- 另外你的SQL语句中,
datetime.today().strftime('%Y-%m-%d')是在循环外计算的,只会取一次当天日期(不过这不是当前问题的核心)。
psycopg2 场景:
- psycopg2 默认关闭自动提交(
autocommit=False),每次创建连接后执行查询会隐式开启事务,若未显式提交/关闭事务,新连接的查询依然无法获取最新提交的数据;同时你每次循环创建新连接却未关闭,可能导致连接泄漏,进一步影响数据可见性。
- psycopg2 默认关闭自动提交(
解决办法
针对 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
相关产品推荐
相关产品推荐

