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

如何在Python中高效排序1亿条CSV数据集并最小化内存占用?

低内存排序1亿条CSV记录的Python解决方案

一、原生Python实现外部排序(分块+多路归并)

无需依赖第三方库,完全基于Python内置模块实现,适合严格控制依赖的场景。

1. 拆分并排序小文件

先将大CSV拆分为若干个内存可容纳的小文件,每个小文件单独排序后保存:

import csv
import os

def split_and_sort_large_csv(input_path, temp_dir, chunk_size=1_000_000, sort_key_idx=0):
    os.makedirs(temp_dir, exist_ok=True)
    chunk_num = 0
    
    with open(input_path, 'r', encoding='utf-8') as infile:
        reader = csv.reader(infile)
        header = next(reader)
        chunk = []
        
        for row in reader:
            chunk.append(row)
            if len(chunk) >= chunk_size:
                # 按指定字段排序(需确保排序键为可排序类型,如时间戳转数值/ datetime)
                chunk.sort(key=lambda x: x[sort_key_idx])
                temp_file = os.path.join(temp_dir, f'sorted_chunk_{chunk_num}.csv')
                with open(temp_file, 'w', newline='', encoding='utf-8') as outfile:
                    writer = csv.writer(outfile)
                    writer.writerow(header)
                    writer.writerows(chunk)
                chunk_num += 1
                chunk = []
        
        # 处理剩余的最后一块
        if chunk:
            chunk.sort(key=lambda x: x[sort_key_idx])
            temp_file = os.path.join(temp_dir, f'sorted_chunk_{chunk_num}.csv')
            with open(temp_file, 'w', newline='', encoding='utf-8') as outfile:
                writer = csv.writer(outfile)
                writer.writerow(header)
                writer.writerows(chunk)

# 调用示例:按第0列(时间戳)排序,每块100万条
split_and_sort_large_csv('large_data.csv', 'temp_sorted_chunks', sort_key_idx=0)

注:chunk_size可根据16GB内存调整,建议设为50万-150万条,避免内存溢出。若时间戳为字符串格式,需先转换为datetime或数值类型再排序,防止字典序排序错误。

2. 多路归并排序后的小文件

用heapq.merge实现多路归并,每次仅从每个小文件读取一行,内存占用极低:

import heapq
import csv
import os

def merge_sorted_chunks(temp_dir, output_path, sort_key_idx=0):
    temp_files = [os.path.join(temp_dir, f) for f in os.listdir(temp_dir) if f.startswith('sorted_chunk_')]
    readers = []
    headers = None
    
    for f in temp_files:
        file = open(f, 'r', encoding='utf-8')
        reader = csv.reader(file)
        headers = next(reader)
        # 包装为(排序键,行内容,文件句柄)格式,适配堆排序逻辑
        def gen(reader_obj, file_obj):
            for row in reader_obj:
                yield (row[sort_key_idx], row, file_obj)
        readers.append(gen(reader, file))
    
    # 执行多路归并并写入结果
    with open(output_path, 'w', newline='', encoding='utf-8') as outfile:
        writer = csv.writer(outfile)
        writer.writerow(headers)
        for _, row, file in heapq.merge(*readers):
            writer.writerow(row)
    
    # 关闭所有临时文件句柄
    for gen_obj in readers:
        try:
            # 从生成器中取出最后一个元素以获取文件句柄并关闭
            last_item = list(gen_obj)[-1]
            last_item[2].close()
        except IndexError:
            pass

# 调用示例
merge_sorted_chunks('temp_sorted_chunks', 'final_sorted_data.csv')

二、用第三方库简化操作

若不想手动实现外部排序逻辑,以下库可快速实现低内存排序:

1. Dask DataFrame

Dask会自动拆分数据集为块,并行处理并归并结果,内存占用可控:

import dask.dataframe as dd

# 按256MB分块读取CSV,可根据内存调整块大小
df = dd.read_csv('large_data.csv', blocksize='256MB')
# 按时间戳字段排序(假设字段名为'timestamp')
sorted_df = df.sort_values(by='timestamp')
# 保存为单个CSV文件
sorted_df.to_csv('final_sorted_data_dask.csv', single_file=True)

安装:pip install dask[complete]

2. Vaex

Vaex采用内存映射技术,无需加载整个数据集到内存,直接对磁盘文件操作:

import vaex

# 读取CSV并创建内存映射
df = vaex.read_csv('large_data.csv')
# 按时间戳排序,生成内存映射格式的排序结果
sorted_df = df.sort(by='timestamp')
# 导出为CSV或更高效的Parquet格式
sorted_df.export_csv('final_sorted_data_vaex.csv')

安装:pip install vaex

三、额外优化技巧

  • 只加载必要字段:若仅需按时间戳排序,可先读取时间戳和行号,排序后再按行号提取原记录,大幅降低内存占用。
  • 使用紧凑存储格式:将临时文件保存为Parquet(比CSV小3-5倍),用pyarrow读写,加快排序和归并速度。
  • 并行分块排序:利用四核CPU,用multiprocessing模块并行处理多个分块,缩短总耗时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 21:12:48