从API端点获取csv.gzip文件并入库的流程优化咨询
Python脚本优化:内存处理压缩CSV+批量数据库导入
原问题背景
现有Python脚本从API获取75000行15列的csv.gzip数据,当前流程是将文件写入磁盘、解压后逐行调用update_or_create写入数据库,全程耗时5-10分钟。需要实现:
- 内存中处理数据,避免磁盘IO
- 批量导入数据库
- 其他可行优化
核心优化方案
1. 内存中处理压缩文件
使用io.BytesIO替代磁盘文件,直接在内存中完成base64解码、gzip解压和CSV读取,彻底消除磁盘IO开销。
2. 批量执行数据库操作
逐行调用update_or_create会产生大量数据库查询和写入请求,是性能瓶颈。可以通过先批量查询已存在的记录,再将数据分为「新增」和「更新」两组分别批量处理,大幅减少数据库交互次数。
优化后的完整代码
import io import csv import base64 import gzip from datetime import datetime from decimal import Decimal from django.db import transaction from .models import pricing response = oauth.get(realm) content = ET.fromstring(response.content) coded_string = content.find('.//pricefile') decoded_string = base64.b64decode(coded_string.text) # 内存中解压并读取CSV with gzip.GzipFile(fileobj=io.BytesIO(decoded_string), mode='rb') as gz_file: csv_file = io.TextIOWrapper(gz_file, encoding='utf-8') reader = csv.reader(csv_file) next(reader) # 跳过表头 current_date = datetime.now() batch_data = [] product_ids = [] # 先收集所有数据和产品ID for row in reader: product_id = row[0] product_ids.append(product_id) batch_data.append({ 'product_id': product_id, 'date': current_date, 'average_price': Decimal(row[1] or 0), 'low_price': Decimal(row[2] or 0), 'high_price': Decimal(row[3] or 0), # 补充其他字段... }) # 批量查询已存在的记录 existing_records = pricing.objects.filter( product_id__in=product_ids, date=current_date ).values_list('product_id', flat=True) existing_ids = set(existing_records) # 拆分新增和更新数据 create_list = [] update_list = [] for data in batch_data: if data['product_id'] in existing_ids: update_list.append(data) else: create_list.append(pricing(**data)) # 批量操作,用事务保证原子性 with transaction.atomic(): # 批量新增 if create_list: pricing.objects.bulk_create(create_list, batch_size=1000) # 批量更新(两种可选方式) if update_list: # 方式1:用bulk_update(需先获取对象) update_objs = pricing.objects.filter( product_id__in=[d['product_id'] for d in update_list], date=current_date ) obj_map = {obj.product_id: obj for obj in update_objs} for data in update_list: obj = obj_map[data['product_id']] obj.average_price = data['average_price'] obj.low_price = data['low_price'] obj.high_price = data['high_price'] # 其他字段赋值... pricing.objects.bulk_update(update_objs, ['average_price', 'low_price', 'high_price', ...], batch_size=1000) # 方式2:用Case/When实现单条SQL批量更新(更高效,适合超大量数据) # from django.db.models import Case, When, Value # update_cases = {} # for field in ['average_price', 'low_price', 'high_price']: # when_clauses = [ # When(product_id=data['product_id'], then=Value(data[field])) # for data in update_list # ] # update_cases[field] = Case(*when_clauses) # pricing.objects.filter( # product_id__in=[d['product_id'] for d in update_list], # date=current_date # ).update(**update_cases)
额外优化建议
- 事务包裹:所有数据库操作放在
transaction.atomic()中,保证数据一致性的同时减少提交次数。 - 固定日期值:循环外只调用一次
datetime.now(),避免每行生成不同时间戳(业务允许时)。 - 调整批量大小:根据数据库性能调整
batch_size(建议1000-5000),平衡内存占用和数据库负载。 - 数据库引擎优化:如果使用PostgreSQL,可尝试
copy_from方法直接导入CSV数据,性能远高于ORM批量操作。 - 关闭自动提交:批量操作前临时关闭Django的自动提交,减少事务开销。
- 使用DictReader:若CSV表头与模型字段名对应,用
csv.DictReader可简化代码,减少索引错误。
内容的提问来源于stack exchange,提问作者Ross
相关产品推荐
相关产品推荐

