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

基于SecurityID高效拆分超大交易数据并追加写入的R性能优化求助

问题描述

需要按SecurityID列高效读取并拆分多份单文件7-10GB的十年交易数据,拆分后按SecurityID生成独立CSV文件:若对应SecurityID的CSV已存在则追加新数据,否则新建文件。现有R代码在小数据集上可正常运行,但处理全量数据时速度极慢,32GB内存的Mac设备上RStudio已白屏运行3天仍在写入文件,需优化性能(如向量化、并行处理等)。

现有代码如下:

# Input of market and transaction data
external.dir <- "my external hard drive"
input.dir.marketdata <-  file.path(external.dir,"marketdata")

# create output of market and transaction data
dir.create(file.path(external.dir, "market_data_split"))
out.dir.marketdata <- file.path(external.dir, "market_data_split")
# loading libraries
library(tidyverse)
library(plyr)
library(data.table)
# Define col to split by as global var
col <- "SecurityID"

# Split and write out the market data
files <- list.files(input.dir.marketdata, pattern = '*MarketData.csv$', full.names = TRUE)
files

for (i in seq_along(files)) {
    input.dat <- data.table::fread(files[i],header = TRUE,stringsAsFactors = TRUE)
    sptdf <- split(input.dat, input.dat[[col]])
    outfile <- as.character(unique(names(sptdf)))
    
    for (j in seq_along(outfile)) {
      new_data <- sptdf[[outfile[j]]]
      outfile.name <- paste0(outfile[j], ".csv")
      check.files <- list.files(out.dir.marketdata, pattern = "*.csv")
      
      if (outfile.name %in% check.files) {
        print(paste0(outfile.name, " already exisit!"))
        exisiting.data <- data.table::fread(file.path(out.dir.marketdata,outfile.name),
                                            header =TRUE,stringsAsFactors = TRUE)
        combined_data <- rbind(exisiting.data, new_data)
        data.table::fwrite(combined_data,
                           file = file.path(out.dir.marketdata, outfile.name),
                           row.names = FALSE,
                           quote = FALSE)
        gc()
     } else {
            print(paste0(outfile.name, " is a new data!"))
            data.table::fwrite(new_data,
                               file = file.path(out.dir.marketdata, outfile.name),
                               row.names = FALSE,
                               quote = FALSE)
            gc()
     }
    }
 }
现有代码性能瓶颈分析
  • 一次性加载全量数据:fread直接读取7-10GB文件到内存,后续split操作进一步消耗内存和CPU,极易触发内存交换拖慢速度。
  • 低效的追加逻辑:每次追加都读取整个已有文件,合并后再重新写入,IO开销呈指数级增长(尤其当目标文件已积累大量数据时)。
  • 频繁文件系统查询:内层循环每次调用list.files检查文件存在性,重复IO操作大幅降低效率。
  • 冗余依赖与参数:加载tidyverse和plyr可能与data.table冲突;stringsAsFactors = TRUE会将字符列转为因子,增加内存开销和处理时间。
  • 过度手动垃圾回收:频繁调用gc()打断正常执行流程,反而降低整体性能。
优化策略
  1. 分块读取大文件:避免一次性加载全量数据,用fread的分块参数逐步处理,控制内存占用。
  2. 直接追加数据:利用fwrite的append = TRUE参数,无需读取已有文件即可追加新数据,彻底消除冗余IO。
  3. 缓存已存在文件列表:提前查询一次输出目录的文件列表,后续循环直接查询内存中的列表,减少IO操作。
  4. 精简依赖包:仅保留data.table,避免包冲突和不必要的内存占用。
  5. 关闭因子转换:去掉stringsAsFactors = TRUE,用字符列处理更高效。
  6. 减少手动GC:让R自动处理垃圾回收,仅在内存紧张时按需调用。
  7. 可选并行处理:若磁盘IO不是瓶颈,可并行处理不同输入文件(注意:并行写同一文件会冲突,仅并行处理不同输入文件)。
优化后的代码
# 初始化路径
external.dir <- "my external hard drive"
input.dir.marketdata <- file.path(external.dir, "marketdata")
out.dir.marketdata <- file.path(external.dir, "market_data_split")
dir.create(out.dir.marketdata, showWarnings = FALSE)

# 仅加载必要包
library(data.table)

# 定义拆分列
split_col <- "SecurityID"

# 获取输入文件列表
input_files <- list.files(input.dir.marketdata, pattern = '*MarketData.csv$', full.names = TRUE)

# 提前缓存已存在的输出文件列表(仅查询一次)
existing_out_files <- list.files(out.dir.marketdata, pattern = "\\.csv$", full.names = FALSE)

# 分块处理每个输入文件,分块大小可根据内存调整
chunk_size <- 1e6 
for (file in input_files) {
    cat("Processing file:", file, "\n")
    
    # 获取文件总行数,用于计算分块数
    total_rows <- fread(file, select = 1, nrows = 0, header = TRUE)$NROW
    num_chunks <- ceiling(total_rows / chunk_size)
    
    for (chunk_idx in 1:num_chunks) {
        cat("Processing chunk", chunk_idx, "/", num_chunks, "\n")
        
        # 分块读取数据
        chunk_data <- fread(
            file,
            header = TRUE,
            skip = ifelse(chunk_idx == 1, 0, (chunk_idx - 1) * chunk_size + 1),
            nrows = chunk_size,
            stringsAsFactors = FALSE
        )
        
        # 按SecurityID分组写入/追加数据
        chunk_data[, {
            out_filename <- paste0(.BY[[split_col]], ".csv")
            out_path <- file.path(out.dir.marketdata, out_filename)
            
            # 检查文件是否存在
            file_exists <- out_filename %in% existing_out_files
            
            # 写入/追加数据,仅新建文件时写表头
            fwrite(
                .SD,
                file = out_path,
                row.names = FALSE,
                quote = FALSE,
                append = file_exists,
                col.names = !file_exists
            )
            
            # 如果是新文件,更新缓存的文件列表
            if (!file_exists) {
                existing_out_files <<- c(existing_out_files, out_filename)
            }
        }, by = split_col]
        
        # 内存紧张时可手动触发GC
        # gc(verbose = FALSE)
    }
}

cat("All files processed!\n")
额外优化建议
  • 使用SSD存储:将输入输出文件放在SSD上,机械硬盘的IO速度是这类任务的主要瓶颈之一。
  • 调整分块大小:根据内存情况调整chunk_size,平衡内存占用和分块次数。
  • 并行处理输入文件:如果有多份输入文件,可使用future.apply并行处理不同文件(避免并行处理同一文件的分块,防止写冲突):
    library(future.apply)
    plan(multisession, workers = 4)  # 根据CPU核心数调整
    future_lapply(input_files, function(file) {
        # 此处复制单个文件的分块处理逻辑
        total_rows <- fread(file, select = 1, nrows = 0, header = TRUE)$NROW
        num_chunks <- ceiling(total_rows / chunk_size)
        
        for (chunk_idx in 1:num_chunks) {
            chunk_data <- fread(
                file,
                header = TRUE,
                skip = ifelse(chunk_idx == 1, 0, (chunk_idx - 1) * chunk_size + 1),
                nrows = chunk_size,
                stringsAsFactors = FALSE
            )
            
            chunk_data[, {
                out_filename <- paste0(.BY[[split_col]], ".csv")
                out_path <- file.path(out.dir.marketdata, out_filename)
                file_exists <- out_filename %in% existing_out_files
                
                fwrite(
                    .SD,
                    file = out_path,
                    row.names = FALSE,
                    quote = FALSE,
                    append = file_exists,
                    col.names = !file_exists
                )
                
                if (!file_exists) {
                    existing_out_files <<- c(existing_out_files, out_filename)
                }
            }, by = split_col]
        }
    })
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 06:43:14