Python中高效读取大型CSV文件列表的优化方案咨询
大CSV文件异步读取+POST优化方案
问题背景
现有多个200MB+且持续增大的CSV文件,原代码读取时会将全量数据加载到内存,导致系统卡顿甚至死机,需要优化为流式分块读取,避免内存溢出。原功能为读取CSV行,按(LabebStoreId, catalog_uuid)分组后POST到API,每组最多2条数据。
原代码核心问题
原代码中data = {}会缓存CSV所有行数据,大文件直接把内存占满,这是卡顿/死机的根本原因。
优化后代码
import asyncio import aiohttp import csv import ast import logging from glob import glob logging.basicConfig( filename="post_script.log", format="%(asctime)s %(message)s", filemode="w" ) logger = logging.getLogger() logger.setLevel(logging.DEBUG) csv.field_size_limit(100000000) headers = {"Accept": "*/*", "Content-Type": "application/json"} # 生成单条row的payload数据 def build_row_payload(row): return { "LabebStoreId": row["LabebStoreId"], "catalog_uuid": row["catalog_uuid"], "lang": row["lang"], "cat_0_name": row["cat_0_name"], "cat_1_name": row["cat_1_name"], "cat_2_name": row["cat_2_name"], "cat_3_name": row["cat_3_name"], "catalogname": row["catalogname"], "description": row["description"], "properties": row["properties"], "price": row["price"], "price_before_discount": row["price_before_discount"], "externallink": row["externallink"], "Rating": row["Rating"], "delivery": row["delivery"], "discount": row["discount"], "instock": row["instock"], "images": ast.literal_eval(row["encoded_images"]), } async def post_single_payload(session, store_id, payload, link_info): """封装POST请求逻辑,复用代码""" response = await session.post( f"http://crawlerapi.labeb.com/api/PCCrawler/Crawl?StoreId={store_id}", json=payload, ) logger.debug(f"Posting to {response.url}") for link in link_info: logger.debug(f"\"\"row: {link}\"\"") logger.info(await response.text()) logger.info("-" * 80) async def post_func(session, file, store_id): current_groups = {} with open(file, "r") as f_in: reader = csv.DictReader(f_in) for row in reader: key = (row["LabebStoreId"], row["catalog_uuid"]) # 当前组添加行 current_groups.setdefault(key, []).append(row) # 组内达到2条,立即POST并清空该组 if len(current_groups[key]) == 2: v = current_groups.pop(key) payload = { "row": build_row_payload(v[0]), "nextRow": build_row_payload(v[1]) } await post_single_payload(session, store_id, payload, [v[0]["externallink"], v[1]["externallink"]]) # 处理文件读完后剩余的单条组 for v in current_groups.values(): if len(v) == 1: payload = {"row": build_row_payload(v[0])} await post_single_payload(session, store_id, payload, [v[0]["externallink"]]) async def main(): try: async with aiohttp.ClientSession(headers=headers) as session: logging.info(f"---starting new---") post_tasks = [] files = glob("*.csv") for file in files: store_id = file.split("_")[0] post_tasks.append( asyncio.create_task(post_func(session, file, store_id)) ) await asyncio.gather(*post_tasks) logging.info(f"---finished---") except Exception as e: logger.error(e) if __name__ == "__main__": loop = asyncio.get_event_loop() loop.run_until_complete(main())
关键优化点
- 流式处理,不缓存全量数据:读取一行就分组,组内达到2条立即POST并删除缓存,内存仅保留当前未处理的分组,占用量极低。
- 代码复用:提取
build_row_payload和post_single_payload函数,减少重复代码,降低维护成本。 - 保持异步并发:依然支持多文件异步处理,同时每个文件内部流式读取,兼顾并发效率与内存控制。
内容的提问来源于stack exchange,提问作者X-something
相关产品推荐
相关产品推荐

