如何用Python合并30GB大型CSV文件?解决内存不足难题
解决大CSV按ID合并的内存不足问题
针对30GB CSV文件按ID合并的内存瓶颈,以下是几种实用的解决方案:
1. Polars懒加载+流式处理
Polars的LazyFrame机制专为大数据场景设计,无需全量加载数据即可完成合并,内存占用远低于pandas。
import polars as pl # 定义所有CSV文件路径 csv_files = ["file1.csv", "file2.csv", ..., "filen.csv"] # 懒加载第一个文件作为基础表 base_df = pl.scan_csv(csv_files[0]) # 依次合并其余文件,延迟执行合并逻辑 for file in csv_files[1:]: current_df = pl.scan_csv(file) base_df = base_df.join(current_df, on="ID", how="outer") # 直接流式写入结果,全程不加载完整数据集到内存 base_df.sink_parquet("merged_result.parquet")
注:Parquet格式比CSV更节省磁盘空间,后续读取也更快,优先推荐。
2. 分块哈希分组合并(适合pandas用户)
通过哈希将相同ID分配到同一数据块,分块独立合并,避免全量加载数据:
import pandas as pd import hashlib def get_id_chunk(id_val, num_chunks=10): # 用ID哈希值确定所属块,确保相同ID在同一块 hash_val = int(hashlib.md5(str(id_val).encode()).hexdigest(), 16) return hash_val % num_chunks num_chunks = 10 # 可根据内存调整块数量 # 初始化第一个文件的分块存储 for chunk_id in range(num_chunks): with pd.read_csv("file1.csv", chunksize=10000) as reader: for part in reader: part["chunk"] = part["ID"].apply(lambda x: get_id_chunk(x, num_chunks)) part[part["chunk"] == chunk_id].to_parquet(f"temp_chunk_{chunk_id}.parquet", mode="append") # 处理其余文件,按块合并 for file in csv_files[1:]: with pd.read_csv(file, chunksize=10000) as reader: for part in reader: part["chunk"] = part["ID"].apply(lambda x: get_id_chunk(x, num_chunks)) for chunk_id in range(num_chunks): chunk_part = part[part["chunk"] == chunk_id] # 读取临时块合并后覆盖写入 temp_df = pd.read_parquet(f"temp_chunk_{chunk_id}.parquet") merged_temp = temp_df.merge(chunk_part, on="ID", how="outer") merged_temp.to_parquet(f"temp_chunk_{chunk_id}.parquet", mode="overwrite") # 合并所有临时块得到最终结果 final_df = pd.concat([pd.read_parquet(f"temp_chunk_{i}.parquet") for i in range(num_chunks)]) final_df.drop("chunk", axis=1).to_csv("merged_result.csv", index=False)
3. 轻量级数据库辅助合并(SQLite)
利用SQLite的磁盘存储特性,将数据导入数据库后用SQL完成JOIN,内存占用极低:
import pandas as pd import sqlite3 conn = sqlite3.connect("temp_merge_db.sqlite") # 分块导入第一个文件到基础表 with pd.read_csv("file1.csv", chunksize=10000) as reader: for part in reader: part.to_sql("base_table", conn, if_exists="append", index=False) # 依次导入其余文件并执行JOIN更新基础表 for idx, file in enumerate(csv_files[1:]): table_name = f"temp_table_{idx+1}" with pd.read_csv(file, chunksize=10000) as reader: for part in reader: part.to_sql(table_name, conn, if_exists="append", index=False) # 执行JOIN并替换基础表 cursor = conn.cursor() cursor.execute(f""" CREATE TABLE temp_base AS SELECT b.*, t.* FROM base_table b LEFT JOIN {table_name} t ON b.ID = t.ID """) cursor.execute("DROP TABLE base_table") cursor.execute("ALTER TABLE temp_base RENAME TO base_table") cursor.execute(f"DROP TABLE {table_name}") conn.commit() # 导出最终合并结果 final_df = pd.read_sql("SELECT * FROM base_table", conn) final_df.to_csv("merged_result.csv", index=False) conn.close()
4. 预处理压缩数据体积
合并前先优化每个CSV的数据类型,减少内存占用:
- 将字符串列转为分类类型(Polars用
pl.Categorical,pandas用dtype="category") - 裁剪不必要的列,仅保留ID和需要合并的字段
- 数值类型向下兼容(比如int64转int32、float64转float32,需确保数值范围允许)
示例(Polars指定schema):
import polars as pl # 自定义紧凑schema custom_schema = { "ID": pl.Utf8, "gender": pl.Categorical, "age": pl.Int32, "score": pl.Float32 } # 用指定schema懒加载文件 base_df = pl.scan_csv("file1.csv", schema=custom_schema)
内容的提问来源于stack exchange,提问作者Spartak Aghababyan
相关产品推荐
相关产品推荐

