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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 08:16:02