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

从AMFI网站提取每日基金NAV数据存入MongoDB的代码优化请求

优化AMFI基金每日NAV数据写入MongoDB的方案

核心思路

把低效的逐条插入改成分批次批量写入,通过控制批次大小平衡内存占用与写入效率,同时优化数据处理流程,减少不必要的内存消耗。

具体优化步骤

1. 分批次批量插入MongoDB

MongoDB的insert_many()相比单条insert_one()能大幅减少网络IO开销。关键是设置合理的批次大小(比如1000-5000条/批,根据本地内存调整),避免一次性加载所有数据到内存。

示例代码改造(替换原逐条插入逻辑):

import pymongo
from pymongo import MongoClient

# 初始化MongoDB连接
client = MongoClient('mongodb://localhost:27017/')
db = client['fund_nav_db']
collection = db['daily_nav']

# 改成迭代器模式,流式返回单条NAV记录(不要一次性加载所有数据)
def fetch_amfi_nav_data():
    # 这里替换成你从AMFI网站下载、解析数据的逻辑
    # 比如逐行解析CSV、分页爬取API,处理一条返回一条
    for nav_record in your_data_source:
        yield nav_record

# 批次插入逻辑
batch_size = 2000  # 内存紧张就调小到500-1000
batch = []
for record in fetch_amfi_nav_data():
    batch.append(record)
    if len(batch) >= batch_size:
        collection.insert_many(batch)
        batch.clear()  # 清空批次释放内存
# 处理最后一批不足batch_size的数据
if batch:
    collection.insert_many(batch)

2. 优化数据加载方式

  • 不要一次性下载全量数据到本地文件再读取,改成流式处理:比如爬取时逐行解析CSV/API响应,处理一条就加入批次,避免大文件占用内存。
  • 如果是下载压缩包(比如AMFI的zip格式NAV文件),直接在内存中解压并流式读取,无需写入本地磁盘。

3. MongoDB端优化

  • 无顺序要求时,设置insert_many(batch, ordered=False):某条数据插入失败不影响其他数据写入,提升容错性和速度。
  • 提前创建索引:比如按基金代码、日期创建复合索引,但要在写入前完成,避免写入时的索引维护开销。
  • 调整连接配置:MongoClient的maxPoolSize设为更大值(比如100),提升并发写入能力;追求速度的话可暂时设w=0(无写入确认),生产环境需权衡数据安全性。

4. 内存监控与动态调优

用psutil监控内存占用,动态调整批次大小,避免内存溢出:

import psutil

def get_memory_usage():
    return psutil.virtual_memory().percent

base_batch_size = 2000
batch = []
for record in fetch_amfi_nav_data():
    batch.append(record)
    # 内存占用超80%或达到批次上限就写入
    if len(batch) >= base_batch_size or get_memory_usage() > 80:
        collection.insert_many(batch)
        batch.clear()
if batch:
    collection.insert_many(batch)
  • 只保留必要字段:解析NAV时丢弃无关冗余字段,减少单条记录的内存占用。

5. 并行处理(可选)

机器有多核的话,用线程池并行爬取数据(注意不要触发AMFI反爬),再统一收集到批次写入:

from concurrent.futures import ThreadPoolExecutor

def process_single_fund(fund_id):
    # 单只基金的NAV数据爬取、解析逻辑
    return nav_records

# 并行爬取,结果合并到批次
with ThreadPoolExecutor(max_workers=5) as executor:
    for records in executor.map(process_single_fund, fund_id_list):
        batch.extend(records)
        if len(batch) >= batch_size:
            collection.insert_many(batch)
            batch.clear()

效果预期

按批次大小2000计算,写入效率至少提升10-100倍,原本5天的任务可缩短至几小时,同时内存占用始终控制在合理范围。

内容的提问来源于stack exchange,提问作者dewashya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 00:25:15