使用asyncio搭配pandas.to_sql()时内存耗尽问题求助
解决异步API获取+SQL保存的内存不足问题
你的判断完全正确:原代码会先把所有API返回的结果全部缓存到内存中,再批量保存,当数据量较大时必然导致内存溢出。核心解决思路是边获取数据边处理保存,避免一次性持有所有数据。以下是具体修改方案:
关键修改点
- 使用
asyncio.as_completed()替代asyncio.gather(),这样每个API任务完成后立刻处理结果,无需等待所有任务结束 - 维护滚动批次,每积累到指定大小就启动保存任务,及时释放内存
- 跟踪所有保存任务,确保最终所有数据都被写入,且等待所有保存操作完成
修改后的完整代码
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
相关产品推荐
相关产品推荐

