【问题标题】:Parallel data.table operations并行数据表操作
【发布时间】:2021-03-11 06:10:16
【问题描述】:

我的脚本中有一个 data.table 对象,我需要将单独的日期和时间列转换为 POSIXct 对象。我正在使用fastPOSIXct() 转换为POSIXct。我发现data.table 操作存在一些瓶颈来执行此操作。主要问题是在转换为POSIXctPOSIXct 转换本身之前连接字符串。我正在使用stri_c() 更快地执行paste0()。有没有办法并行化这个计算来加快它?正在读取的 csv 文件是一个大约 2gb 的大文件。

data.table structure

  index = match(file, csv_files)
  print(paste0("Starting File #", index, ". Running Active Future Filter."))
  csv1 = fread(file = file)
  # create row for date in Date format
  csv1[, date := as.Date(integer())]
  csv1[, trade_date := as.character(trade_date)]
  csv1[, date := as.Date(trade_date[1], "%Y%m%d"), by = trade_date]
  csv1[, active_exp := character()]
  csv1[, active_exp := get.active.future(date = date[1]), by = date]
  nrow(csv1)
  csv1 = csv1[csv1$active_exp == as.character(csv1$contract_delivery_date),]
  # Add POSIXct
  # keep only essential columns
  csv1 = csv1[,-c(3, 6, 9, 23, 11:21)]
  print(paste0("Running POSIX"))
  # establish date-time POSIX for tick
  csv1[, date_time_char := as.character()]
  # creates date time character type to be converted to POSIXct
  csv1[, date_time_char := stri_c(date[1], trade_time[1], sep = " "), by = list(date, trade_time)]
  
  # convert date time character type to POSIXct
  csv1[, date_time := fastPOSIXct(date_time_char[1]), by = date_time_char]
  # time shift for date when time is between 17:00 and 00:00
  print("Making time adjustment")
  csv1[,adj_date_time := time.shift(date_time[1]), by = date_time]
  # convert date
  csv1[, date := as.Date(adj_date_time, tz = "UTC")]
  # overwrite current loaded file

【问题讨论】:

  • 欢迎您!请使您的示例可重现。另外,关于这个问题,我不确定这是“data.table 操作的一些瓶颈”,因为它从字符串转换为POSIXct 很慢。
  • @Cole 为无法包含数据而道歉,它是付费的,不能在公共论坛上发布。我可以发送 str() 或任何其他信息。
  • 可复制并不意味着共享私人数据。这意味着花时间提供可以模拟问题的数据集。
  • @chinsoon12 我也有同样的想法,但我认为这可能是一个性能问题。如果您有很多重复的日期,那么按字符串分组然后对一个字符串而不是所有字符串进行转换可能会更快。

标签: r data.table large-data posixct


【解决方案1】:

这应该会有所帮助(未经测试,因为该示例不可重现):

index = match(file, csv_files)

csv1 = fread(file = file)
# create row for date in Date format
csv1[, trade_date := as.character(trade_date)]
csv1[,
     c("date", "active_exp") := {
       date1 = as.Date(.BY[[1L]], "%Y%m%d")
       active_exp = get.active.future(date1)
       .(date1, active_exp)},
     by = trade_date]

csv1 = csv1[active_exp == contract_delivery_date), -c(3, 6, 9, 23, 11:21)]

# creates date time character type to be converted to POSIXct
csv1[, 
     c("date_time", "adj_date_time", "date") = {
       date_time = fastPOSIXct(stri_c(.BY[[1L]]), .BY[[2L]], sep = " ")
       adj_date_time = time.shift(date_time)
       date = as.Date(adj_date_time, tz = "UTC")
     }
     , by = .(date, trade_time)
]

有多个地方对或多或少相同的数据进行了分组。每次我们进行分组时,这意味着我们正在调用forder() 并且必须对数据进行分组。相反,我们可以尝试一次做所有事情。

为了使这更快,我建议只freading 您需要的列。相关,导入文件后,我会做过滤器active_exp == contract_delivery_date

编辑 您可以查看的一件事是使用IDateITime 创建POSIXct 列。最慢的部分可能是stri_c / paste,然后是分组或fasttime::fastPOSIXct。这是避免粘贴的方法,尽管它依赖于一些 Rcpp 来帮助解析时间戳列。

Rcpp::cppFunction("
IntegerVector to_time(std::vector< std::string > x) {
  //for format hh:mm:ss 
  int n = x.size();
  IntegerVector out(n);
  
  for (int i = 0; i < n; i++){;
    const std::string xi = x[i]; 
    const int hour = stoi(xi.substr(0, 2));
    const int minute = stoi(xi.substr(3, 2));
    const int second = stoi(xi.substr(6, 2));
    out[i] = hour * 60 * 60 + minute * 60 + second;
  }
  return(out);
}
")

csv1[,
     date_time := {
       date1 = as.IDate(as.character(trade_date), "%Y%m%d")
       time = as.ITime(to_time(trade_time))
       .(as.POSIXct(date1, time, tz = "UTC"))
     }]

当复制小数据集一百万次时,这需要 3 秒,而fasttime::fastPOSIXct() 需要 20 秒。将在下面发布基准测试。

您还应该考虑删除原始帖子中的by = 部分。如果有很多重复的日期时间,那么是的,按组分组有时会更有效率。但如果有很多独特的日期时间,最好跳过分组步骤。

csv1 = csv1[rep(seq_len(.N), 1e6L)]

## for use in different use cases
date_col = as.Date(as.character(csv1$trade_date), "%Y%m%d")
IDate_col = as.IDate(date_col)
hour_col = as.ITime(to_time(csv1$trade_time))
bench::mark(
 as_char_to_date = as.Date(as.character(csv1$trade_date), "%Y%m%d")
 ,
 as.ITime(to_time(csv1$trade_time))
 ,
 as_date_to_idate = as.IDate(date_col)
 , as.POSIXct(IDate_col, hour_col, tz = "UTC")
 ,  use_fasttime = fasttime::fastPOSIXct(paste(date_col, csv1$trade_time), "UTC")
 , check = FALSE
)

## # A tibble: 5 x 13
##   expression                                       ## min   median `itr/sec`
##   <bch:expr>                                  <bch:tm> <bch:tm>     <dbl>
## 1 as_char_to_date                                3.03s    3.03s    0.330 
## 2 as.ITime(to_time(csv1$trade_time))             2.03s    2.03s    0.494 
## 3 as_date_to_idate                              13.2ms   13.5ms   23.1   
## 4 as.POSIXct(IDate_col, hour_col, tz = "UTC") 929.14ms 929.14ms    1.08  
## 5 use_fasttime                                  20.39s   20.39s    0.0490

【讨论】:

  • 这大大加快了进程,但仍然需要(不合理的 imo)很长时间。
  • 查看编辑。还要注意Rcpp 可以相对容易地与OpenMP 并行。我用 8 个线程对其进行了测试,在 ITime 步骤中,从 Date -> POSIXct 开始大约 1.8 秒,而 Date -> fasttime 则为 18 秒。
  • 这很好用。将括号date1 = as.IDate(as.character(trade_date, "%Y%m%d")) 之一编辑为date1 = as.IDate(as.character(trade_date), "%Y%m%d")
  • 已编辑。你可以随意接受。另外,我会删除您的答案,尽管可能会编辑您的问题以包含数据集。
【解决方案2】:

编辑: 正如@Cole 所建议的那样,我使用了这个 C++ 函数以及 IDate 和 ITime 来加快速度。我没有对此进行基准测试,但可以说它的运行速度比以前的代码快 20 倍。代码如下。

for(file in csv_files){

  index = match(file, csv_files)
  csv1 = fread(file = file, drop = c(3, 6, 9, 23, 11:21))
  
  options(warn = -1)
  # get date and time and classify as POSIXct
  csv1[,
       c("date_time", "date", "time") := {
         
         time = as.ITime(to_time(trade_time))
         date1 = as.IDate(as.character(trade_date), "%Y%m%d")
         if(hour(time) >= 17){
            POSIX = as.POSIXct(date1, time, tz = "UTC") - days(1)
            .(POSIX, as.Date(POSIX), time)
         }else{
            POSIX = as.POSIXct(date1, time, tz = "UTC")
            .(POSIX, as.Date(POSIX), time)
         }
       }]
  options(warn = 0)
  # covert trade date to character format
  csv1[, trade_date := as.character(trade_date)]
  # run active_exp
  csv1[,
       active_exp := get.active.future(as.Date(date[1])),
       by = date]
  # filter only active future
  csv1 = csv1[active_exp == as.character(contract_delivery_date)]
  # write file
  file_path = paste0("C:/Users/ocean/Documents/CME Data/Active Daily Data CSVs/", index, ".csv")
  fwrite(csv1, file = file_path)
  
  print(paste0("Progress = ", round(index/length(csv_files), 3)*100, "%"))
}

旧: 我使用了@Cole 提供的代码,但仍然运行了很长时间。我稍微编辑了代码以修复一些语法,如下所示。 文件 = csv_files[1] 索引 = 匹配(文件,csv_files)

csv1 = fread(file = file, drop = c(3, 6, 9, 23, 11:21))
# create row for date in Date format
csv1[, trade_date := as.character(trade_date)]
csv1[,
     c("date", "active_exp") := {
       date1 = as.Date(.BY[[1L]], "%Y%m%d")
       active_exp = get.active.future(date1)
       .(date1, active_exp)},
     by = trade_date]

csv1 = csv1[active_exp == contract_delivery_date]

# creates date time character type to be converted to POSIXct
csv1[, 
     c("date_time", "adj_date_time", "date") := {
       date_time = fastPOSIXct(stri_c(.BY[[1L]], .BY[[2L]], sep = " "))
       adj_date_time = time.shift(date_time)
       date = as.Date(adj_date_time, tz = "UTC")
     }
     , by = .(date, trade_time)
]

为了提供一个可重复的示例,数据的dput() 如下。

    structure(list(trade_date = c(20200115L, 20200115L, 20200115L, 
20200115L, 20200115L, 20200115L), trade_time = c("17:00:00", 
"17:00:00", "17:00:00", "17:00:00", "17:00:00", "17:00:00"), 
    trade_sequence_number = c(9028350L, 9028357L, 9028366L, 9028394L, 
    9028397L, 9028400L), session_indicator = c("E", "E", "E", 
    "E", "E", "E"), ticker_symbol = c("ES", "ES", "ES", "ES", 
    "ES", "ES"), future_option_index_indicator = c("F", "F", 
    "F", "F", "F", "F"), contract_delivery_date = c(2003L, 2003L, 
    2003L, 2003L, 2003L, 2003L), trade_quantity = c(176L, 0L, 
    3L, 2L, 4L, 10L), strike_price = c(0L, 0L, 0L, 0L, 0L, 0L
    ), trade_price = c(3287.75, 3287.75, 3288, 3288, 3288, 3288
    ), ask_bid_type = c(NA, NA, NA, NA, NA, NA), indicative_quote_type = c(NA, 
    NA, NA, NA, NA, NA), market_quote = c(NA, NA, NA, NA, NA, 
    NA), close_open_type = c("", "O", "", "", "", ""), valid_open_exception = c(NA, 
    NA, NA, NA, NA, NA), post_close = c(NA, NA, NA, NA, NA, NA
    ), cancel_code_type = c(NA, NA, NA, NA, NA, NA), insert_code_type = c(NA, 
    NA, NA, NA, NA, NA), fast_late_indicator = c(NA, NA, NA, 
    NA, NA, NA), cabinet_indicator = c(NA, NA, NA, NA, NA, NA
    ), book_indicator = c(NA, NA, NA, NA, NA, NA), entry_date = c(20200114L, 
    20200114L, 20200114L, 20200114L, 20200114L, 20200114L), exchange_code = c("XCME", 
    "XCME", "XCME", "XCME", "XCME", "XCME")), row.names = c(NA, 
-6L), class = c("data.table", "data.frame"), .internal.selfref = <pointer: 0x000001ef853c1ef0>)

考虑到文件的大小(38 个文件为 2-4gb),我是只需要学习 C++ 还是让它通宵运行?

【讨论】:

  • 你能控制csv中提供的数据吗? fread 将采用 yyyy-mm-dd 格式的日期并自动将其读取为 IDate 格式。
  • @Cole 我无法控制使用 cURL API 请求从数据库中提取的数据
猜你喜欢
  • 2019-08-12
  • 2014-11-14
  • 2017-10-27
  • 2018-03-31
  • 2020-03-23
  • 2021-10-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多