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

如何在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)) %&gt;% 
  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) %&gt;% 
    filter(updated_dt == max(updated_dt),
           PAIRING_STATUS == "A") %&gt;% 
  select(PAIRING_POSITION, CREW_ID, PAIRING_NO, PAIRING_DATE, CREW_INDICATOR, updated_dt) %&gt;% 
  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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 14:50:03