如何在R中并行化独立代码块?以数据库查询与清洗为例
在R中并行化独立数据库查询与数据清洗任务
能否在R中并行化独立代码块内的所有操作?例如我有3个分别查询数据库不同表的代码块,后续还会进行数据清洗。能否将每个代码块拆分为独立进程,在不同核心上同时运行,让进程1、2、3并行执行而非按顺序运行?
进程1
q_flighthistory <- "SELECT * FROM CT_FLIGHT_HISTORY;" # WHERE FLIGHT_DATE BETWEEN '2024-07-01' AND '2024-07-02';" dbGetQuery(q_flighthistory, "view_flighthistory") # Unique flight history for that date clean_flighthistory <- view_flighthistory %>% filter(SEGMENT_STATUS == "A") %>% select(c(2:10), (12:13), (55:59),EQUIPMENT) %>% #filter(!duplicated(UNIQUE_ID)) %>% rename(SCHED_DEPARTURE_TIME = SCHED_DEPARTURE_TIME_RAW)%>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep= " ")) %>% group_by(FLIGHT_NO, FLIGHT_DATE, DEPARTING_CITY, ARRIVAL_CITY, SCHED_DEPARTURE_TIME) %>% filter(updated_dt == max(updated_dt)) %>% mutate(temp_id = cur_group_id()) %>% filter(!duplicated(temp_id)) %>% select(!c(12:16)) %>% select(!c(5:7)) %>% select(!c(10:11)) rm(view_flighthistory) gc()
进程2
q_flightleg <- "SELECT * FROM CT_FLIGHT_LEG;" # WHERE PAIRING_DATE BETWEEN '2024-06-28' AND '2024-07-02';" dbGetQuery(q_flightleg, "view_flightleg") flight_leg_join <- view_flightleg %>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep= " "), shed_dept_hour = hour(SCHED_DEPARTURE_TIME), DEADHEAD = if_else(is.na(DEADHEAD), "C", DEADHEAD)) %>% relocate(updated_dt, .after=PAIRING_NO) %>% group_by(PAIRING_POSITION, DEPARTING_CITY, ARRIVAL_CITY, SCHED_DEPARTURE_DATE, shed_dept_hour, PAIRING_NO, PAIRING_DATE) %>% filter(DEADHEAD == "C", updated_dt == max(updated_dt)) %>% mutate(temp_id = cur_group_id(), .after = CREW_INDICATOR) %>% filter(!duplicated(temp_id)) %>% ungroup() %>% select(PAIRING_NO, PAIRING_DATE, FLIGHT_NO, SCHED_DEPARTURE_DATE, SCHED_ARRIVAL_DATE, DEPARTING_CITY, ARRIVAL_CITY, PAIRING_POSITION) rm(view_flightleg) gc()
进程3
# MasterPairing Query q_masterpairing <- "SELECT * FROM CT_MASTER_PAIRING;" # WHERE PAIRING_DATE BETWEEN '2024-06-28' AND '2024-07-02';" dbGetQuery(q_masterpairing, "view_masterpairing") join_masterpairing <- view_masterpairing %>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep = " ")) %>% relocate(updated_dt, .after=PAIRING_NO) %>% group_by(PAIRING_DATE, CREW_ID) %>% filter(updated_dt == max(updated_dt), PAIRING_STATUS == "A") %>% select(PAIRING_POSITION, CREW_ID, PAIRING_NO, PAIRING_DATE, CREW_INDICATOR, updated_dt) %>% ungroup() rm(view_masterpairing) gc()
实现并行化的方法
你的三个任务完全独立,非常适合并行执行。推荐使用future+furrr组合,语法简洁且兼容dplyr操作,以下是具体实现步骤:
1. 安装并加载依赖包
install.packages(c("future", "furrr", "dplyr", "lubridate")) library(future) library(furrr) library(dplyr) library(lubridate)
2. 设置并行策略
使用multisession模式创建独立子进程,指定使用的核心数(建议留1个核心给系统):
plan(multisession, workers = availableCores() - 1)
3. 将每个进程封装为独立函数
注意:并行子进程无法共享主进程的数据库连接,必须在函数内部单独建立连接:
进程1函数
process_flighthistory <- function() { # 这里添加你的数据库连接代码,例如: # con <- dbConnect(RPostgres::Postgres(), dbname = "your_db", host = "your_host", user = "user", password = "pwd") q_flighthistory <- "SELECT * FROM CT_FLIGHT_HISTORY;" view_flighthistory <- dbGetQuery(con, q_flighthistory) # 注意参数顺序:连接在前,查询在后 clean_flighthistory <- view_flighthistory %>% filter(SEGMENT_STATUS == "A") %>% select(c(2:10), (12:13), (55:59), EQUIPMENT) %>% rename(SCHED_DEPARTURE_TIME = SCHED_DEPARTURE_TIME_RAW) %>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep= " ")) %>% group_by(FLIGHT_NO, FLIGHT_DATE, DEPARTING_CITY, ARRIVAL_CITY, SCHED_DEPARTURE_TIME) %>% filter(updated_dt == max(updated_dt)) %>% mutate(temp_id = cur_group_id()) %>% filter(!duplicated(temp_id)) %>% select(!c(12:16)) %>% select(!c(5:7)) %>% select(!c(10:11)) rm(view_flighthistory) gc() dbDisconnect(con) # 关闭连接释放资源 return(clean_flighthistory) }
进程2函数
process_flightleg <- function() { # 建立数据库连接 # con <- dbConnect(...) q_flightleg <- "SELECT * FROM CT_FLIGHT_LEG;" view_flightleg <- dbGetQuery(con, q_flightleg) flight_leg_join <- view_flightleg %>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep= " "), shed_dept_hour = hour(SCHED_DEPARTURE_TIME), DEADHEAD = if_else(is.na(DEADHEAD), "C", DEADHEAD)) %>% relocate(updated_dt, .after=PAIRING_NO) %>% group_by(PAIRING_POSITION, DEPARTING_CITY, ARRIVAL_CITY, SCHED_DEPARTURE_DATE, shed_dept_hour, PAIRING_NO, PAIRING_DATE) %>% filter(DEADHEAD == "C", updated_dt == max(updated_dt)) %>% mutate(temp_id = cur_group_id(), .after = CREW_INDICATOR) %>% filter(!duplicated(temp_id)) %>% ungroup() %>% select(PAIRING_NO, PAIRING_DATE, FLIGHT_NO, SCHED_DEPARTURE_DATE, SCHED_ARRIVAL_DATE, DEPARTING_CITY, ARRIVAL_CITY, PAIRING_POSITION) rm(view_flightleg) gc() dbDisconnect(con) return(flight_leg_join) }
进程3函数
process_masterpairing <- function() { # 建立数据库连接 # con <- dbConnect(...) q_masterpairing <- "SELECT * FROM CT_MASTER_PAIRING;" view_masterpairing <- dbGetQuery(con, q_masterpairing) join_masterpairing <- view_masterpairing %>% mutate(updated_dt = paste(UPDATE_DATE, UPDATE_TIME, sep = " ")) %>% relocate(updated_dt, .after=PAIRING_NO) %>% group_by(PAIRING_DATE, CREW_ID) %>% filter(updated_dt == max(updated_dt), PAIRING_STATUS == "A") %>% select(PAIRING_POSITION, CREW_ID, PAIRING_NO, PAIRING_DATE, CREW_INDICATOR, updated_dt) %>% ungroup() rm(view_masterpairing) gc() dbDisconnect(con) return(join_masterpairing) }
4. 并行执行所有任务
# 任务列表 tasks <- list(process_flighthistory, process_flightleg, process_masterpairing) # 并行执行,结果以列表返回 results <- future_map(tasks, ~.x()) # 提取各自的结果 clean_flighthistory <- results[[1]] flight_leg_join <- results[[2]] join_masterpairing <- results[[3]]
关键注意事项
- 数据库连接:每个子进程必须单独建立/关闭连接,避免连接冲突或资源泄漏。
- 内存管理:子进程的内存独立,函数内的
rm()和gc()操作仅作用于当前子进程。 - 核心数控制:不要占满所有CPU核心,留1个给系统进程,防止系统卡顿。
- 任务独立性:确保任务之间没有数据依赖(你的场景完全满足),否则无法并行。
内容的提问来源于stack exchange,提问作者Eizy
相关产品推荐
相关产品推荐

