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

8GB内存下合并500GB级金融CSV数据的高效处理方案咨询

一、首先排查代码低级错误

你给出的示例代码存在变量名不匹配问题:读入的表命名为table1、table2、weight,合并时调用的是data1、data2、data3,先修正这个问题再排查内存问题。


二、不同规模场景的解决方案

方案1:最小改动快速修复(临时处理少量大文件)

仅依赖pandas即可实现,无需额外搭建环境:

  • 读CSV时强制指定字段类型,大幅降低内存占用:
    • 分类字段(fund_ticker、TICKER)指定为category类型,比object类型省80%以上内存
    • 日期字段Date指定parse_dates参数直接解析为datetime类型,不要存为字符串
    • 数值字段根据精度要求选用float32/int32等小长度类型,不要默认用float64/int64
  • 先将所有CSV分块转成Parquet格式:单文件5G的CSV转成Parquet后一般不到1G,列式存储特性可以让你合并时只加载需要的关联键和字段,内存占用骤降
    示例代码:
    import pandas as pd
    # 分块转CSV到Parquet
    chunk_size = 10**6
    for i, chunk in enumerate(pd.read_csv(r'C:\data\data1.csv', 
                                          usecols=lambda col: col not in ["Version"],
                                          dtype={"fund_ticker": "category", "TICKER": "category"},
                                          parse_dates=["Date"],
                                          chunksize=chunk_size)):
        chunk.to_parquet(f'data1_part{i}.parquet', index=False)
    # 后续同逻辑转另外两个文件,再读取Parquet文件合并即可
    
  • 额外排查:确认每个表的关联键['fund_ticker', 'TICKER', 'Date']是否存在大量重复值,重复值会导致合并时产生笛卡尔积,内存占用指数级上涨,合并前先按需去重。

方案2:1TB以内数据稳定处理方案(嵌入式SQL)

无需搭建数据库服务,用Python自带的sqlite3即可实现磁盘级运算,内存占用稳定在1G以内:

  • 步骤:
    1. 新建本地sqlite数据库文件
    2. 将三个CSV分块导入到sqlite的三个表中,给每个表的关联键创建联合索引
    3. 写SQL语句执行三表JOIN,结果直接导出到目标文件
      示例代码:
    import sqlite3
    import pandas as pd
    
    conn = sqlite3.connect('fund_data.db')
    chunk_size = 10**6
    
    # 导入第一个表
    for chunk in pd.read_csv(r'C:\data\data1.csv', 
                            usecols=lambda col: col not in ["Version"],
                            chunksize=chunk_size):
        chunk.to_sql('table1', conn, if_exists='append', index=False)
    # 同逻辑导入table2、weight表
    # 创建联合索引加速查询
    conn.execute('CREATE INDEX idx1 ON table1(fund_ticker, TICKER, Date)')
    conn.execute('CREATE INDEX idx2 ON table2(fund_ticker, TICKER, Date)')
    conn.execute('CREATE INDEX idx3 ON weight(fund_ticker, TICKER, Date)')
    conn.commit()
    
    # 执行JOIN并导出结果
    join_sql = '''
    SELECT * FROM table1
    INNER JOIN table2 USING(fund_ticker, TICKER, Date)
    INNER JOIN weight USING(fund_ticker, TICKER, Date)
    '''
    # 分块读取结果避免内存溢出
    for i, chunk in enumerate(pd.read_sql(join_sql, conn, chunksize=chunk_size)):
        chunk.to_parquet(f'merged_result_part{i}.parquet', index=False)
    
    conn.close()
    

方案3:长期TB级金融数据处理方案

适合高频处理大规模金融数据的场景,搭建后长期使用效率提升明显:

  • 存储层:所有原始数据统一转成Parquet格式,按Date、fund_ticker等高频查询字段分区存储,比CSV存储成本降低70%以上,查询效率提升10倍以上
  • 计算层:本地用Dask框架即可,API和pandas几乎完全兼容,自动实现分块运算、内存调度,8G内存即可轻松处理TB级数据;后续数据量增长到10TB以上时可以无缝迁移到Spark集群。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 06:48:03