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

调用3GB PostgreSQL表为何占用内存远超表本身?

解决psycopg2拉取大表内存暴涨问题及异步POST优化

问题核心原因

内存暴涨的关键在于cur.fetchall()会一次性将3GB的表数据全部加载到内存,再转换为Pandas DataFrame时,Python对象(元组、DataFrame列结构)的内存开销远大于原始数据库存储,直接导致内存占用飙升至27GB。同步数据库读取与异步POST请求的不协调也会加剧内存积压。


解决方案

1. 使用服务器端游标分批读取数据

通过服务器端游标让数据库分批返回数据,避免一次性将所有数据传输到客户端,从根源上减少内存占用。

修改后的分批查询函数:

def query_database_batch(
    query: str,
    batch_size: int = 1000,
    query_args: str = None,
    dbname: str = os.environ.get('DBNAME'),
    user: str = os.environ.get('DBUSER'),
    password: str = os.environ.get('DBPASSWORD'),
    host: str = os.environ.get('LOCAL_BIND_ADDRESS'),
    port: int = int(os.environ.get('LOCAL_BIND_PORT'))
):
    with psycopg2.connect(
            dbname=dbname,
            user=user,
            password=password,
            host=host,
            port=port
        ) as conn:
        # 创建服务器端游标(需指定名称开启该模式)
        with conn.cursor(name='large_table_cursor') as cur:
            args = [query, query_args] if query_args else [query]
            cur.execute(*args)
            columns = [desc[0] for desc in cur.description] if cur.description else []
            
            # 循环分批获取数据
            while True:
                records = cur.fetchmany(batch_size)
                if not records:
                    break
                # 每批生成小DataFrame返回,或直接处理原始数据
                yield pd.DataFrame(data=records, columns=columns)

2. 结合asyncio实现异步POST发送

由于psycopg2是同步库,需通过asyncio.to_thread将数据库读取操作放到线程中,避免阻塞异步事件循环,同时用aiohttp实现异步POST请求。

示例代码:

import asyncio
import aiohttp

async def send_batch(session: aiohttp.ClientSession, batch_df):
    # 将批次数据转换为API接收的格式(如JSON)
    data = batch_df.to_dict('records')
    async with session.post('https://目标接口地址.com', json=data) as resp:
        if resp.status != 200:
            print(f"批次发送失败,状态码: {resp.status}")

async def main():
    query = sql.SQL("SELECT * from table")
    async with aiohttp.ClientSession() as session:
        # 遍历分批读取的数据
        for batch_df in query_database_batch(query):
            await send_batch(session, batch_df)
            # 主动释放当前批次内存
            del batch_df

if __name__ == "__main__":
    asyncio.run(main())

3. 优化DataFrame内存占用

转换DataFrame时指定更高效的数据类型,进一步降低内存开销:

# 针对不同列定义内存友好的数据类型
dtype_mapping = {
    '整数列': 'int32',
    '字符串列': 'category',
    '浮点数列': 'float32'
}

# 在生成DataFrame时传入类型映射
yield pd.DataFrame(data=records, columns=columns, dtype=dtype_mapping)

额外建议

  • 如果业务不需要DataFrame,可以直接处理原始元组数据,跳过DataFrame转换环节,进一步减少内存占用。
  • 根据本地内存和接口性能调整batch_size(如500、2000),找到内存消耗与传输效率的平衡点。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 07:10:06