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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 22:50:03