【问题标题】:Merge millions of csv files with different headers合并数百万个具有不同标题的 csv 文件
【发布时间】:2019-12-19 22:27:10
【问题描述】:

我有数百万个带有不同标题的 csv 文件,我想将它们合并到一个大数据框中。

我的问题是我尝试过的解决方案有效,但是太慢了!顺便说一句,我可以在我的实验室中使用 Sparklyr 在多节点集群上工作,这个大数据工具会有帮助吗?

文件如下所示:

文件1

标头1,标头3,标头5

a,b,c

文件2

标头4,标头2

e,f

文件3

header2,header6

a,c

我想把它们合并成:

header1,header2,header3,header4,header5,header6

a,,b,,c, ,f,,e,, ,a,,,,c

我曾尝试将它们直接与 R 绑定,但程序在服务器中运行几天后崩溃。代码如下所示:

library(plyr)
library(dplyr)
library(readr)


csvfiles <- list.files(pattern = "file\\d+.csv") 

for (i in 1:length(csvfiles)) {
  assign(paste0("files", i),read_csv(csvfiles[[i]]))
}

csvlist <- mget(ls(pattern = "files\\d"))

result <- data.frame()

for (i in 1:length(csvlist)){
  my_list <- list(result,csvlist[[i]])
  result <- rbindlist(my_list,use.names=TRUE, fill=TRUE)
 }

然后我尝试首先使用sedawkcsvtk等命令行工具提取标题。我使用的代码如下所示

for file in $(ls file*.csv); do cat $file | sed "2 d" | csvtk transpose   >> name_combined.csv; done

awk  '{ lines[$1] = $0 } END { for (l in lines) print lines[l] }' name_combined.csv >> long_head.csv

我得到名为 long_head.csv 的 csv 文件,它看起来像这样(实际上我得到了超过 3000 列)

header1,header2,header3,header4,header5,header6

然后我在dplyr 中使用bind_rows。我想首先输出具有相同列的多个 csv 文件,然后将它们全部合并。

library(readr)
library(dplyr)

csvfiles <- list.files(pattern = "file\\d+.csv")
long_head <- read_csv("long_head.csv")

new_file <- paste("new_file",1:length(csvfiles),sep = "")

for (i in 1:length(csvfiles)) {
         bind_rows(long_head,read_csv(csvfiles[[i]]))  %>% 
            write_csv (file = paste0(new_file [[i]], ".csv"))
}

代码一天只能输出大约 10 万个 csv 文件,这意味着我必须等待整整一个月才能让这些 csv 文件合并它们。

我也试过不写多个csv文件直接组合起来:

library(readr)
library(dplyr)

csvfiles <- list.files(pattern = "file\\d+.csv")

long_head <- read_csv("long_head.csv")


for (i in 1:length(csvfiles)) {
  a <- bind_rows(read_csv(csvfiles[[i]]),long_head)
  result <- rbind(a,long_head)
}

它跑得更快,但也落后于我的预期。

【问题讨论】:

  • 试试data.table::rbindlist()
  • 嗨,你自己尝试了什么?你提到你有一些运行速度很慢的代码。你能也给我们看看吗?
  • 首先,转置所有文件。所以每一行都变成了列。然后开始一个接一个地追加文件。之后,转置你的大文件,所以每一列都变成了行。瞧,您已经根据需要合并了文件
  • 我已经明确了我的问题...我的问题是我无法一次将所有 csv 文件读入 R,因此速度很慢。

标签: r csv awk bigdata


【解决方案1】:

您可以在下面找到一种使用 GNU awk 的方法,它将完全读取文件。它将执行以下操作:

  • 读取每个文件的头部,并关闭文件
  • 如果发现新的标题元素,请将其添加到当前已知元素的末尾。例如,存在以下标头:

    file1: A,B,D
    file2: A,C,E
    file3: A,E,D
    

    输出标题

    A,B,D,C,E
    
  • 分析完所有头文件后,它会完整读取所有文件并在需要时用空字段重写整个文件。

这个脚本使用What's the most robust way to efficiently parse CSV using awk?

创建一个文件merge_csv.awk,内容如下:

BEGIN {
   OFS=","
   FPAT="[^,]*|\042[^\042]+\042"
   # keep track of the original argument count
   argc_start=ARGC
}

# Read header and process
# header names are stored as array index in the array "header"
# header order is stored in the array header_order
#    header_order[field_index] = header_name
(FNR == 1) && (ARGIND < argc_start) {
    for(i=1;i<=NF;++i) if (!($i in header)) { header[$i]; header_order[++nf_out]=$i } 
    # add file to end of argument list to be reprocessed
    ARGV[ARGC++] = FILENAME
    # process the next file
    nextfile
}

# Print headers in output file
(FNR == 1) && (ARGIND == argc_start) {
    for(i=1;i<=nf_out;++i) printf header_order[i] (i==nf_out ? ORS : OFS)
}

# Use array h to keep track of the column_name and corresponding field_index
# h[column_name] = field_index
(FNR == 1) { delete h; for(i=1;i<=NF;++i) h[$i]=i; next }

# print record
{
    # process all fields
    for(i=1;i<=nf_out;++i) {
        # get field index using h
        j = h[header_order[i]]+0
        # if field index is zero, print empty field
        printf (j == 0 ? "" : $j) (i==nf_out ? ORS : OFS)
    }
}

现在你可以运行脚本了

$ awk -f merge_csv.awk *.csv > output.csv

这不适用于大量 CSV 文件。这可以通过以下方式解决。假设您有一个文件filelist.txt,其中包含您想要的所有文件(可以由find 生成),然后将上述脚本添加为:

BEGIN {
   OFS=","
   FPAT="[^,]*|\042[^\042]+\042"
}

# Read original filelist, and build argument list
(FNR == NR) { ARGV[ARGC++]=$0; argc_start=ARGC; next }

# Read header and process
# header names are stored as array index in the array "header"
# header order is stored in the array header_order
#    header_order[field_index] = header_name
(FNR == 1) && (ARGIND < argc_start) {
    for(i=1;i<=NF;++i) if (!($i in header)) { header[$i]; header_order[++nf_out]=$i } 
    # add file to end of argument list to be reprocessed
    ARGV[ARGC++] = FILENAME
    # process the next file
    nextfile
}

# Print headers in output file
(FNR == 1) && (ARGIND == argc_start) {
    for(i=1;i<=nf_out;++i) printf header_order[i] (i==nf_out ? ORS : OFS)
}

# Use array h to keep track of the column_name and corresponding field_index
# h[column_name] = field_index
(FNR == 1) { delete h; for(i=1;i<=NF;++i) h[$i]=i; next }

# print record
{
    # process all fields
    for(i=1;i<=nf_out;++i) {
        # get field index using h
        j = h[header_order[i]]+0
        # if field index is zero, print empty field
        printf (j == 0 ? "" : $j) (i==nf_out ? ORS : OFS)
    }
}

现在您可以将代码运行为:

$ awk -f merge_csv.awk filelist.txt

如果您的文件列表确实太大,您可能需要使用split 并使用循环创建各种临时 CSV 文件,这些文件可以在第二次甚至第三次运行时再次合并。

【讨论】:

  • 不错的答案。如果 OP 真的有数百万个文件,我担心awk ... *.csv 和 ARG_MAX。
  • @MarkSetchell,你是对的。我添加了一个更新。
【解决方案2】:
  • 使用dir 和模式来选择文件名;
  • 添加源文件列,以后会很有用;
  • 更简单的循环调用;
  • 强制所有列为字符,读取多个文件时最安全的选择,readr 解析猜测功能将在遇到字段不匹配时中止。

注意:运行 16 个文件的测试总是使我的计算机崩溃,介于 15MB 771 列 census.csv 和 180MB 1.6M 行 beer_reviews.csv 之间。

library(readr)
library(dplyr)

setwd("/home/username/R/csv_test")

csvfiles <- dir(pattern = "\\.csv$")

csvdata  <- tibble(filename=c("Source File"))

for (i in csvfiles) {
  tmpfile <- read_csv(i, col_types = cols(.default = "c"))
  tmpfile$filename <- i
  csvdata <- bind_rows(csvdata, tmpfile)
}
csvdata
# A tibble: 1,622,379 x 874

...

定时 10 个文件测试运行,总共 20k 行和 100 列。在 R 中:

 user  system elapsed 
0.678   0.008   0.685 

还有这个页面上的 awk 脚本:

real    0m2.202s
user    0m2.175s
sys     0m0.025s

【讨论】:

    【解决方案3】:

    这是一个具有挑战性的问题,需要同时考虑速度和内存消耗。

    如果我理解正确,OP 想要合并数百万个小 csv 文件。根据样本数据,每个文件仅包含 2 行:第一行的标题和第二行的字符数据。列数和列名可能因文件而异。但是,所有列都是相同的数据类型字符。

    OP 的第一次尝试以及M. Viking's answer 都在迭代地增长结果对象。这是非常低效的,因为它需要一遍又一遍地复制相同的数据。此外,两者都使用readr 包中的read_csv(),这也不是最快的csv阅读器。

    为了避免迭代地增长结果对象,所有文件都被读入一个列表,然后使用rbindlist() 一次性组合。然后将最终结果存储为 csv 文件:

    library(data.table)
    file_names <- list.files(pattern = "file\\d+.csv")
    result <- rbindlist(lapply(file_names, fread), use.names=TRUE, fill=TRUE)
    fwrite(result, "result.csv")
    

    从 OP 的预期结果看来,列应该按列名排序。这可以通过

    library(magrittr)
    setcolorder(result, names(result) %>% sort())
    

    通过引用重新排列 data.table 对象的列,即不复制整个对象。

    性能

    现在,让我们看看处理时间。对于基准测试,我创建了 100k 个文件(请参阅下面的 数据 部分),这与 OP 的目标数量相去甚远,但可以得出结论。

    在我的电脑上,整个处理时间大约是 5 分钟:

    bench::workout({
      fn <- list.files(pattern = "file\\d+.csv")
      tmp_list <- lapply(fn, function(x) fread(file = x, sep =",", header = TRUE, colClasses = "character") )
      result <- rbindlist(tmp_list, use.names=TRUE, fill=TRUE)
      setcolorder(result, names(result) %>% sort())
      fwrite(result, "result.csv")
    }, 1:5)
    
    # A tibble: 5 x 3                                                                                                                    
      exprs       process     real
      <bch:expr> <bch:tm> <bch:tm>
    1 1           562.5ms 577.19ms
    2 2             1.81m    4.52m
    3 3            14.05s   15.55s
    4 4           15.62ms  175.1ms
    5 5              2.2s    7.72s
    

    在这里,我使用了 bench 包中的 workout() 函数对单个表达式进行计时,以便识别花费最多时间的语句。批量读取 csv 文件。

    对象的大小也很重要。合并后的data.table result 100k 行1000 列占用800 MB,临时列表仅占11%。这是因为有许多空单元格。

    pryr::object_size(result)
    
    800 MB
    
    pryr::object_size(tmp_list)
    
    87.3 MB
    

    顺便说一下,结果文件"result.csv" 在磁盘上的大小为 98 MB。

    结论

    计算时间似乎不是主要问题,而是存储结果所需的内存。

    如果读取 100k 个文件大约需要 5 分钟,我猜读取 1M 个文件可能需要大约 50 分钟。

    但是,对于具有 3000 列的 1M 文件,生成的 data.table 可能需要 10 * 3 = 30 倍的内存,即 24 GB。临时列表可能只需要大约 900 MB。因此,重新考虑生成的数据结构可能是值得的。

    对不同的文件阅读器进行基准测试

    以上时序表明,超过 90% 的计算时间用于读取数据文件。因此,对读取 CSV 文件的不同方法进行基准测试是合适的:

    • read.csv() 来自基础 R
    • read_csv() 来自 readr
    • fread() 来自 data.table

    为了方便用户,所有三个函数都提供猜测文件的某些特征,例如字段分隔符或数据类型。这可能需要额外的计算时间。因此,这些函数也使用明确声明的文件参数进行基准测试。

    对于基准测试,使用了bench 包,因为它还测量分配的内存,这可能是计算时间之外的另一个限制因素。对不同数量的文件重复基准测试,以研究对内存消耗的影响。

    library(data.table)
    library(readr)
    file_names <- list.files(pattern = "file\\d+.csv")
    bm <- press(
      n_files = c(1000, 2000, 5000, 10000),
      {
        fn <- file_names[seq_len(n_files)] 
        mark(
          fread = lapply(fn, fread),
          fread_p = lapply(fn, function(x) fread(file = x, sep =",", header = TRUE, colClasses = "character")),
          # fread_pp = lapply(fn, fread, sep =",", header = TRUE, colClasses = "character"),
          read.csv = lapply(fn, read.csv),
          read.csv_p = lapply(fn, read.csv, colClasses = "character"),
          read_csv = lapply(fn, read_csv),
          read_csv_p = lapply(fn, read_csv, col_types = cols(.default = col_character())),
          check = FALSE,
          min_time = 10
        )
      }
    )
    

    结果可视化

    library(ggplot2)
    ggplot(bm) + aes(n_files, median, color = names(expression)) + 
      geom_point() + geom_line() + scale_x_log10()
    ggsave("median.png")
    ggplot(bm) + aes(n_files, mem_alloc, color = names(expression)) + 
      geom_point() + geom_line() + scale_x_log10()
    ggsave("mem_alloc.png")
    ggplot(bm) + aes(median, mem_alloc, color = names(expression)) + geom_point() + 
      facet_wrap(vars(n_files))
    ggsave("mem_allov_vs_median.png")
    

    在比较中值执行时间时,我们可以观察到(请注意双对数刻度)

    • 计算时间几乎随着文件数量线性增加;
    • 显式传递文件参数(计时被命名为..._p)总是比猜测参数带来性能提升,尤其是read_csv()
    • 对于读取许多小文件,read.csv()fread() 快,而 read_csv() 则慢一些。

    在比较分配的内存时,我们可以观察到(再次注意双对数刻度)

    • 内存消耗几乎随文件数量线性增加;
    • 显式传递文件参数(时序命名为..._p)对内存分配没有显着影响,但read_csv() 除外,其中猜测参数的成本似乎相当高;
    • fread() 分配的内存明显少于其他两个读取器。

    如上所述,速度内存消耗在这里可能至关重要。现在,read.csv() 似乎是速度方面的最佳选择,而fread() 是内存消耗方面的最佳选择,从下面的多面散点图中可以看出。

    我个人的选择是更喜欢fread()(更少的内存消耗)而不是read.csv()(更快),因为我的电脑上的内存是有限的,不能轻易扩展。您的里程可能会有所不同。

    数据

    以下代码用于创建 100k 示例文件:

    library(magrittr)   # piping used to improve readability
    n_files <- 10^4L
    max_header <- 10^3L
    avg_cols <- 4L
    headers <- sprintf("header%02i", seq_len(max_header))
    set.seed(1L)   # to ensure reproducible results
    for (i in seq_len(n_files)) {
      n_cols <- rpois(1L, avg_cols - 1L) + 1L # exclude 0
      header <- sample(headers, n_cols)
      file_name <- sprintf("file%i.csv", i)
      sample(letters, n_cols, TRUE) %>% 
        as.list() %>% 
        as.data.frame() %>%
        set_names(header) %>% 
        data.table::fwrite(file_name)
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-03-20
      • 1970-01-01
      • 2020-03-29
      • 2012-12-08
      • 2017-10-05
      • 2021-06-15
      • 2019-07-18
      相关资源
      最近更新 更多