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

使用asyncio搭配pandas.to_sql()时内存耗尽问题求助

解决异步API获取+SQL保存的内存不足问题

你的判断完全正确:原代码会先把所有API返回的结果全部缓存到内存中,再批量保存,当数据量较大时必然导致内存溢出。核心解决思路是边获取数据边处理保存,避免一次性持有所有数据。以下是具体修改方案:

关键修改点

  1. 使用asyncio.as_completed()替代asyncio.gather(),这样每个API任务完成后立刻处理结果,无需等待所有任务结束
  2. 维护滚动批次,每积累到指定大小就启动保存任务,及时释放内存
  3. 跟踪所有保存任务,确保最终所有数据都被写入,且等待所有保存操作完成

修改后的完整代码

import asyncio
import polars as pl
import pandas as pd
import itertools


async def save_to_land(transactions, table_name, conn_object, semaphore):
    def process_and_save(transactions):
        df = pl.DataFrame(transactions, infer_schema_length=2000)
        # 这里保留你的Polars处理逻辑
        # ...
        # 注意:如果conn_object是共享连接,需确保线程安全,建议使用连接池
        df.to_pandas().to_sql(table_name, conn_object, if_exists='append', index=False)

    if transactions:
        async with semaphore:
            await asyncio.to_thread(process_and_save, transactions)


async def get_transactions(api, practice_id, year, incremental, with_deleted, date_modified, semaphore):
    async with semaphore:
        transactions = await api.get_transactions(
            practice_id, year, incremental=incremental, 
            with_deleted=with_deleted, date_modified=date_modified
        )
        await asyncio.sleep(0.5)  # 保留API限流逻辑
        return transactions


async def get_transactions_and_save_to_land(
    api, practice_ids, table_name, conn_object, 
    start_year, end_year, incremental, with_deleted, date_modified, 
    batch_size=100000
):
    # API请求并发控制
    sem_api = asyncio.Semaphore(4)
    # 数据库写入并发控制
    sem_db = asyncio.Semaphore(3)

    if not incremental:
        # 创建所有API任务,但不立即等待全部完成
        get_tasks = [
            asyncio.create_task(
                get_transactions(api, practice_id, year, incremental, with_deleted, date_modified, sem_api)
            )
            for practice_id, year in itertools.product(practice_ids, range(start_year, end_year + 1))
        ]

        batch_transactions = []
        save_tasks = []

        # 迭代完成的API任务,实时处理结果
        for completed_task in asyncio.as_completed(get_tasks):
            transactions = await completed_task
            if not transactions:
                continue
            
            batch_transactions.extend(transactions)
            
            # 批次达到阈值,启动保存任务
            if len(batch_transactions) >= batch_size:
                # 复制当前批次,避免后续修改影响保存任务
                save_batch = batch_transactions.copy()
                save_task = asyncio.create_task(
                    save_to_land(save_batch, table_name, conn_object, sem_db)
                )
                save_tasks.append(save_task)
                # 清空批次,准备下一轮积累
                batch_transactions = []

        # 处理剩余的未达批次阈值的数据
        if batch_transactions:
            save_task = asyncio.create_task(
                save_to_land(batch_transactions, table_name, conn_object, sem_db)
            )
            save_tasks.append(save_task)

        # 等待所有保存任务完成
        await asyncio.gather(*save_tasks)

额外优化建议

  • 避免Pandas转换开销:Polars支持直接写入SQL数据库(pl.write_database()),无需转成Pandas DataFrame,能减少内存占用和转换时间。示例:
    # 替换原有的df.to_pandas().to_sql(...)
    df.write_database(
        table_name=table_name,
        connection=conn_object,
        if_exists="append",
        engine="adbc"  # 需安装对应的ADBC驱动,如adbc-sqlserver-driver
    )
    
  • 连接线程安全:如果你的conn_object是单个数据库连接,在多线程环境下(asyncio.to_thread())可能出现线程安全问题,建议使用SQLAlchemy连接池,让每个保存任务获取独立连接。
  • 调整批次大小:根据内存情况和数据库写入性能,微调batch_size参数,平衡内存占用和写入效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 13:10:16