在Julia中按指定数量CSV文件分批创建DataFrame及内存优化咨询
Julia实现CSV文件分批合并为DataFrame及大文件处理方案
一、分批合并CSV到DataFrame的实现步骤
首先确保安装所需依赖包:
using Pkg Pkg.add(["CSV", "DataFrames", "Filesystem"])
核心代码实现:
using CSV, DataFrames, Filesystem # 目标目录路径(raw字符串避免转义问题) dir_path = raw"C:\Users\me\Files" # 获取目录下所有CSV文件并按名称排序(保证批次顺序稳定) csv_files = sort(filter(file -> endswith(file, ".csv"), readdir(dir_path, join=true))) # 设置每批文件数量 batch_size = 5 # 用字典存储各批次DataFrame(比单独命名变量更灵活) batch_dfs = Dict{String, DataFrame}() # 分批处理并合并 for (batch_num, file_batch) in enumerate(Iterators.partition(csv_files, batch_size)) # 合并当前批次的所有CSV文件 merged_df = vcat([CSV.read(file, DataFrame) for file in file_batch]...) # 存入字典,键格式为df_1、df_2... batch_dfs["df_$batch_num"] = merged_df end # 示例:访问第一个批次的DataFrame # first_batch_df = batch_dfs["df_1"]
代码说明
readdir(dir_path, join=true)返回带完整路径的文件名,避免路径拼接错误sort()确保文件顺序固定,保证"前5个、接下来5个"的逻辑符合预期Iterators.partition()自动将文件列表按指定批次大小拆分vcat()合并同批次的多个DataFrame,需确保各CSV文件的列结构一致
二、内存不足时的处理方案
如果不确定内存是否足够,可根据数据规模选择以下方案:
1. 流式读取CSV(减少单文件内存占用)
使用CSV.jl的CSV.Rows迭代行,或指定chunksize分块读取:
# 分块读取单个大CSV并处理 for chunk in CSV.Rows(dir_path * "\\large_file.csv", chunk_size=10000) # 对当前块进行处理,无需加载整个文件到内存 chunk_df = DataFrame(chunk) # ... 自定义处理逻辑 end
2. 使用数据库存储与查询
将CSV导入SQLite数据库,通过SQL查询分批获取数据,避免全量加载:
using SQLite, CSV # 创建本地SQLite数据库 db = SQLite.DB("csv_data.db") # 逐个导入CSV到数据库表(表名为CSV文件名去掉后缀) for file in csv_files table_name = replace(basename(file), ".csv" => "") CSV.write(db, table_name, CSV.read(file, DataFrame)) end # 示例:合并前5个CSV对应的表到DataFrame df_1 = DBInterface.execute(db, """ SELECT * FROM File1 UNION ALL SELECT * FROM ABC UNION ALL SELECT * FROM x7b_320 -- 继续添加剩余2个表名 """) |> DataFrame # 若单个表过大,用LIMIT/OFFSET分批查询 df_chunk = DBInterface.execute(db, "SELECT * FROM large_table LIMIT 10000 OFFSET 0") |> DataFrame
3. 采用内存映射格式存储
将数据转换为Arrow格式,支持内存映射,无需全量加载到内存:
using Arrow # 合并批次后导出为Arrow文件 for (batch_num, df) in batch_dfs Arrow.write("batch_$batch_num.arrow", df) end # 内存映射读取Arrow文件,仅加载需要的部分 mapped_df = Arrow.Table("batch_1.arrow") |> DataFrame
4. 分块处理后立即导出
如果不需要同时保留所有批次的DataFrame,处理完一批就导出为文件,释放内存:
for (batch_num, file_batch) in enumerate(Iterators.partition(csv_files, batch_size)) merged_df = vcat([CSV.read(file, DataFrame) for file in file_batch]...) # 导出为CSV或Arrow文件 CSV.write("batch_$batch_num.csv", merged_df) # 无需保留merged_df,内存自动回收 end
内容的提问来源于stack exchange,提问作者Coco
相关产品推荐
相关产品推荐

