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

如何基于输入文件并行化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.

Practical Parallel Solution for Processing 7k+ Parquet Files by Player ID

Core Problem Breakdown

You’ve got two main hurdles:

  1. Slow sequential processing: 7k+ files with 1.5M rows each can’t be handled in a loop.
  2. 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: arrow is far faster than base R for reading Parquet, and readr is 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 the str_extract regex (e.g., use ^\\d+ for numeric IDs, or ^[^_]+ to capture everything before the first underscore).
  • Temp storage: For HPC, consider setting temp_directory to a high-speed scratch partition instead of the system temp dir to speed up I/O.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 19:22:48