如何在Python中高效处理与过滤大型CSV文件?
处理GB级CSV的内存高效过滤与计算方案
一、优化Pandas分块读取的实际操作
用chunksize后内存仍高,通常是没提前裁剪数据或优化数据类型,调整方向如下:
- 指定低内存数据类型:读取时通过
dtype参数将高精度类型转成低内存替代项,比如float64转float32、重复率高的字符串用category类型:
import pandas as pd dtype_spec = { "id": "int32", "category_col": "category", "value": "float32" } chunk_iter = pd.read_csv("large_data.csv", dtype=dtype_spec, chunksize=10_000, usecols=["id", "category_col", "value"]) filtered_chunks = [] for chunk in chunk_iter: # 提前过滤,只保留符合条件的行,减少后续处理量 filtered_chunk = chunk[(chunk["value"] > 100) & (chunk["category_col"] == "A")] # 对过滤后的块执行计算 filtered_chunk["calculated_col"] = filtered_chunk["value"] * 1.5 filtered_chunks.append(filtered_chunk) # 手动释放内存 del chunk final_result = pd.concat(filtered_chunks, ignore_index=True)
- 按需加载列:用
usecols只读取需要处理的列,彻底避免加载无关数据,这是降内存最直接的手段。
二、Dask DataFrame的最优过滤实现
Dask天生适配大文件并行处理,过滤逻辑和Pandas几乎一致,无需大幅改动代码:
import dask.dataframe as dd # 读取CSV,Dask自动分块,可通过blocksize自定义块大小(如"100MB") ddf = dd.read_csv("large_data.csv", dtype=dtype_spec, usecols=["id", "category_col", "value"]) # 基础过滤 filtered_ddf = ddf[(ddf["value"] > 100) & (ddf["category_col"] == "A")] # 新增计算列 filtered_ddf["calculated_col"] = filtered_ddf["value"] * 1.5 # 结果集可放入内存则转Pandas,否则继续用Dask处理 final_result = filtered_ddf.compute()
- 复杂逻辑处理:若过滤涉及自定义函数,用
map_partitions将函数应用到每个分块,自动并行执行:
def complex_filter(chunk): mask = chunk["text_col"].str.contains(r"target_pattern", regex=True) & (chunk["value"] > 50) return chunk[mask] filtered_ddf = ddf.map_partitions(complex_filter)
三、SQLite的轻量级数据库方案
SQLite无需额外服务,导入CSV后用SQL做过滤计算,内存占用极低,适合多轮查询场景:
- 导入CSV到SQLite:
import sqlite3 import csv conn = sqlite3.connect("temp_data.db") cursor = conn.cursor() # 提前定义表结构,优化列类型 cursor.execute(""" CREATE TABLE IF NOT EXISTS data ( id INTEGER, category_col TEXT, value REAL ) """) # 批量导入比逐行插入快 with open("large_data.csv", "r") as f: reader = csv.reader(f) next(reader) # 跳过表头 cursor.executemany("INSERT INTO data VALUES (?, ?, ?)", reader) conn.commit()
- 执行过滤与计算:
# 用SQL完成过滤和计算,结果转Pandas DataFrame query = """ SELECT id, category_col, value, value * 1.5 AS calculated_col FROM data WHERE value > 100 AND category_col = 'A' """ final_result = pd.read_sql(query, conn) # 给过滤列建索引,提升后续查询速度 cursor.execute("CREATE INDEX idx_value_category ON data(value, category_col)") conn.commit() conn.close()
四、轻量级纯Python方案
若不想依赖第三方库,可用csv模块结合多进程加速逐行处理:
import csv from multiprocessing import Pool def process_row(row): if float(row["value"]) > 100 and row["category_col"] == "A": row["calculated_col"] = str(float(row["value"]) * 1.5) return row return None def read_csv_rows(file_path): with open(file_path, "r") as f: reader = csv.DictReader(f) for row in reader: yield row if __name__ == "__main__": with Pool(processes=4) as pool: results = pool.map(process_row, read_csv_rows("large_data.csv")) final_results = [res for res in results if res is not None]
总结建议
- 熟悉Pandas的话,优先优化分块参数(
dtype、usecols)+ 提前过滤,快速解决问题; - 复杂计算或需并行,直接用Dask,学习成本低,性能提升明显;
- 多轮查询或复杂过滤逻辑,用SQLite,利用SQL的查询优化能力;
- 纯Python环境追求轻量,用
csv+多进程组合。
内容的提问来源于stack exchange,提问作者KirillovaVika
相关产品推荐
相关产品推荐

