使用arrow::open_dataset读取CSV时skip_rows不生效及分区提取问题
解决Arrow读取大体积CSV的Schema错误与自定义分区问题
1. 修复Schema不匹配错误
你的错误根源是Arrow默认从文件开头推断Schema,而前2行是单列元数据,导致推断出仅1列,但后续数据实际是296列数值网格。解决方法是手动指定Schema,跳过元数据行并强制使用正确的列结构:
- 首先定义296个浮点型列的Schema(适配十进制数值数据):
library(arrow) library(tidyverse) # 创建296个float64类型的列Schema csv_schema <- schema( !!!set_names(map(1:296, ~field(type = float64())), paste0("col", 1:296)) )
- 在
open_dataset中指定该Schema,同时确保跳过前2行,并且关闭自动Schema推断的干扰:
csv_data <- open_dataset( "path", format = "csv", schema = csv_schema, skip_rows = 2, skip_rows_inference = TRUE, # 推断Schema时也跳过前2行 col_names = paste0("col", 1:296) # 给无列名数据指定列名 )
2. 从自定义文件名提取分区(date和time)
针对my_file_Time[yyyymmddThhmmss].csv格式的文件名,使用Arrow的自定义分区提取函数,通过字符串位置截取日期和时间:
# 自定义分区提取逻辑 custom_partitioning <- partitioning( extract = function(file_path) { file_name <- basename(file_path) # 截取日期(位置16-23:yyyymmdd) date_str <- substr(file_name, 16, 23) # 截取时间(位置25-30:hhmmss) time_str <- substr(file_name, 25, 30) list( date = as.Date(date_str, "%Y%m%d"), time = as_hms(strptime(time_str, "%H%M%S")) # 转为time32类型 ) }, schema = schema(date = date32(), time = time32("s")) # 明确分区字段类型 )
将此分区规则传入open_dataset即可完成自定义分区配置。
3. 完整代码与内存优化
整合所有步骤,写入Parquet格式(Arrow的分块读取机制不会导致内存溢出):
library(tidyverse) library(arrow) library(hms) # 创建目标目录 dir.create("converted_parquet", recursive = TRUE, showWarnings = FALSE) # 定义CSV Schema csv_schema <- schema( !!!set_names(map(1:296, ~field(type = float64())), paste0("col", 1:296)) ) # 自定义分区规则 custom_partitioning <- partitioning( extract = function(file_path) { file_name <- basename(file_path) date_str <- substr(file_name, 16, 23) time_str <- substr(file_name, 25, 30) list( date = as.Date(date_str, "%Y%m%d"), time = as_hms(strptime(time_str, "%H%M%S")) ) }, schema = schema(date = date32(), time = time32("s")) ) # 读取CSV数据集 csv_data <- open_dataset( "path", format = "csv", schema = csv_schema, skip_rows = 2, skip_rows_inference = TRUE, col_names = paste0("col", 1:296), partitioning = custom_partitioning ) # 写入Parquet分区数据集 write_dataset( csv_data, "converted_parquet", format = "parquet", partitioning = c("date", "time"), max_rows_per_file = 100000 # 控制单个Parquet文件大小,进一步优化内存 )
关键说明
- 手动Schema:完全规避了Arrow自动推断时的列数不匹配问题,确保读取296列数值数据。
- 自定义分区:通过字符串截取从非Hive风格文件名中提取分区字段,无需重命名文件。
- 内存友好:Arrow的
open_dataset采用分块读取机制,不会将整个文件加载到内存,完美解决fread的内存不足问题。
内容的提问来源于stack exchange,提问作者billyleaf
相关产品推荐
相关产品推荐

