如何基于输入文件并行化R脚本以高效完成大量Parquet文件的玩家ID拆分导出任务
Hey there! Let's work through your problem—you're dealing with a classic large-scale data processing challenge, and we can fix both the speed and parallel writing conflict issues with a practical, runnable solution tailored for HPC environments.
Core Problem Breakdown
You’ve got two main hurdles:
- Slow sequential processing: 7k+ files with 1.5M rows each can’t be handled in a loop.
- Parallel write conflicts: Multiple processes appending to the same CSV will cause corruption, even with
append = TRUE.
We’ll solve this with a two-step approach: parallel processing of input files (using modern, easy-to-use parallel tools) and safe temporary file handling to avoid write collisions.
Required Packages
First, install and load the tools we’ll use—these are more intuitive than doParallel and work reliably in HPC:
install.packages(c("furrr", "future", "arrow", "dplyr", "purrr", "readr", "stringr")) library(furrr) library(future) library(arrow) library(dplyr) library(purrr) library(readr) library(stringr)
Step 1: Define the File Processing Function
This function reads a single Parquet file, splits it by player ID, and writes each player’s data to a unique temporary file (no append needed here—each temp file is exclusive to one process/player segment):
process_parquet <- function(file_path, temp_dir) { # Read the Parquet file (arrow is fast for this!) df <- arrow::read_parquet(file_path) # Split the data frame by player ID split_player_data <- df %>% dplyr::group_split(ID, .keep = FALSE) # Write each player's segment to a unique temp file purrr::walk(split_player_data, function(player_segment) { player_id <- unique(player_segment$ID) # Create a unique temp filename to avoid conflicts (uses process ID + random number) temp_filename <- file.path(temp_dir, paste0(player_id, "_", Sys.getpid(), "_", sample.int(1e6, 1), ".csv")) # Write the segment (no append needed—this is a new file every time) readr::write_csv(player_segment, temp_filename, col_names = FALSE) }) }
Step 2: Define the Merge Function
After processing all files, we’ll merge all temporary files for each player into a single final CSV. This step is sequential (safe for writing) and ensures no data corruption:
merge_temp_files <- function(temp_dir, output_dir) { # Get all temp CSV files all_temp_files <- list.files(temp_dir, pattern = "\\.csv$", full.names = TRUE) # Group temp files by player ID (adjust the regex if your ID format differs!) player_file_groups <- all_temp_files %>% basename() %>% stringr::str_extract("^PL\\d+") %>% # Matches IDs like PL1, PL2—update if your IDs are different purrr::set_names(all_temp_files) %>% split(., .) # Merge each player's temp files into one final CSV purrr::walk(player_file_groups, function(files_for_player) { player_id <- unique(stringr::str_extract(basename(files_for_player), "^PL\\d+")) final_output_path <- file.path(output_dir, paste0(player_id, ".csv")) # Handle column names: only write them once if the file doesn't exist if (!file.exists(final_output_path)) { # Read the first temp file to get column names first_segment <- readr::read_csv(files_for_player[1], col_names = TRUE) readr::write_csv(first_segment, final_output_path, col_names = TRUE) # Skip the first file for appending files_to_append <- files_for_player[-1] } else { files_to_append <- files_for_player } # Append all remaining temp files purrr::walk(files_to_append, function(temp_file) { segment_data <- readr::read_csv(temp_file, col_names = FALSE) readr::write_csv(segment_data, final_output_path, append = TRUE, col_names = FALSE) }) }) }
Step 3: Full Runnable Pipeline
This ties everything together, with parallelization optimized for HPC:
# -------------------------- # Set your paths here! # -------------------------- input_directory <- "Input_Folder_Path" # Parent directory with Parquet files output_directory <- "Output_Folder_Path" # Where final player CSVs will go temp_directory <- file.path(tempdir(), "player_temp_data") # Temp storage (uses system temp or set your own) # Create directories if they don't exist dir.create(output_directory, recursive = TRUE, showWarnings = FALSE) dir.create(temp_directory, recursive = TRUE, showWarnings = FALSE) # Get all Parquet file paths all_parquet_files <- list.files( path = input_directory, pattern = "*.snappy.parquet$", recursive = TRUE, full.names = TRUE ) # Set up parallel processing for HPC: # - Use multicore for Linux/macOS (faster) # - Use multisession for Windows # - Automatically uses all available CPU cores plan(multicore, workers = availableCores()) # Swap to plan(multisession) if on Windows # Process all files in parallel future_map(all_parquet_files, process_parquet, temp_dir = temp_directory) # Merge temp files into final player CSVs merge_temp_files(temp_directory, output_directory) # Clean up temporary files (optional but recommended) unlink(temp_directory, recursive = TRUE)
Key Optimizations & Notes
- No write conflicts: By using unique temporary files per process/player segment, we avoid the risk of parallel processes overwriting the same CSV.
- Speed:
arrowis far faster than base R for reading Parquet, andreadris optimized for CSV operations. - HPC-friendly: The parallel plan automatically uses all available CPU cores—no need to manually specify resources.
- ID format adjustment: If your player IDs don’t follow the
PL\d+pattern, update thestr_extractregex (e.g., use^\\d+for numeric IDs, or^[^_]+to capture everything before the first underscore). - Temp storage: For HPC, consider setting
temp_directoryto a high-speed scratch partition instead of the system temp dir to speed up I/O.
内容的提问来源于stack exchange,提问作者November2Juliet

