20GB大CSV逐行处理HTTP请求的Python脚本优化咨询
大CSV文件异步处理脚本的速度优化建议
我有一个20GB、约1亿行的CSV文件需要处理,该CSV结构简单,包含var1、var2、var3、timestamp(yyyy-MM-dd)四列,所有列均为字符串类型。需求为逐行读取CSV,将日期转换为指定格式后构造请求,逐行发送至服务器(不支持批量发送数据)。我已编写如下Python脚本,希望获取关于该脚本的速度优化建议:
import pandas as pd import string from datetime import datetime import aiohttp import asyncio url_data = 'http://localhost:8080/test' template = string.Template(('''<?xml version="1.0" encoding="UTF-8"?> <request> <const1>CONST1</const1> <const2>CONST2</const2> <var1>${VAR1}</var1> <timestamp>${DATETIME}</timestamp> <const3>CONST3</const3> <var2>${VAR2}</var2> <var3>${VAR3}</var3> </request>''')) async def process(): async with aiohttp.ClientSession() as session: for chunk in pd.read_csv('test.csv', chunksize=50000): await asyncio.gather(asyncio.create_task(send(session, chunk))) async def send(session, chunk): for index, row in chunk.iterrows(): date = datetime.strptime(row['DATETIME'], '%Y-%m-%d') row['DATETIME'] = date.strftime("%Y-%m-%dT%H:%M:%SZ") await session.post(url=url_data, data=template.safe_substitute(row)) if __name__ == '__main__': asyncio.run(process())
核心优化建议
1. 真正发挥异步并发优势
原脚本的异步逻辑完全未生效:send函数内逐行await请求,本质是串行执行。正确做法是将每个请求包装为独立任务批量提交,同时用信号量控制并发数,避免压垮服务器:
- 用
asyncio.Semaphore限制并发请求量(建议50-200,根据服务器承受能力调整) - 批量收集任务后用
asyncio.gather执行,不要逐个等待
2. 替换Pandas为轻量CSV读取工具
Pandas读取CSV会构建DataFrame,带来额外内存和CPU开销,纯逐行处理场景完全没必要。改用标准库的csv.DictReader,读取速度更快,内存占用更低。
3. 砍掉冗余的日期转换
原代码中datetime.strptime解析日期再strftime格式化的操作完全多余——原日期是yyyy-MM-dd格式,直接拼接T00:00:00Z即可,跳过datetime对象转换,能节省大量CPU时间。
4. 优化XML模板渲染
string.Template.safe_substitute性能不如直接使用字符串format或f-string。预定义模板字符串,用format方法填充变量,渲染速度会显著提升。
5. 控制任务批量大小
一次性将1亿个任务加入事件循环会导致内存爆炸,可每积累N个任务(比如1000个)就执行一次gather,然后清空任务列表,避免内存过载。
优化后的示例代码
import csv import asyncio import aiohttp url_data = 'http://localhost:8080/test' # 预定义XML模板,用字符串format提升渲染速度 XML_TEMPLATE = '''<?xml version="1.0" encoding="UTF-8"?> <request> <const1>CONST1</const1> <const2>CONST2</const2> <var1>{var1}</var1> <timestamp>{timestamp}</timestamp> <const3>CONST3</const3> <var2>{var2}</var2> <var3>{var3}</var3> </request>''' # 限制并发请求数,根据服务器性能调整 MAX_CONCURRENT = 50 async def send_single_request(session, semaphore, row): # 直接拼接日期,跳过datetime转换步骤 row['timestamp'] = f"{row['timestamp']}T00:00:00Z" xml_payload = XML_TEMPLATE.format(**row) async with semaphore: async with session.post(url=url_data, data=xml_payload) as resp: # 可根据需求添加响应校验,比如检查状态码 await resp.text() async def process_csv(): semaphore = asyncio.Semaphore(MAX_CONCURRENT) async with aiohttp.ClientSession() as session: pending_tasks = [] with open('test.csv', 'r', encoding='utf-8') as csv_file: reader = csv.DictReader(csv_file) for row in reader: task = asyncio.create_task(send_single_request(session, semaphore, row)) pending_tasks.append(task) # 每积累1000个任务就执行一次,避免内存过载 if len(pending_tasks) >= 1000: await asyncio.gather(*pending_tasks) pending_tasks = [] # 处理剩余的任务 if pending_tasks: await asyncio.gather(*pending_tasks) if __name__ == '__main__': asyncio.run(process_csv())
内容的提问来源于stack exchange,提问作者J.Doe
相关产品推荐
相关产品推荐

