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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 15:25:14