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以内:
- 步骤:
- 新建本地sqlite数据库文件
- 将三个CSV分块导入到sqlite的三个表中,给每个表的关联键创建联合索引
- 写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
相关产品推荐
相关产品推荐

