Debugowanie i równoważenie obciążenia kodu równoległego
Obsługuj błędy w procesach roboczych i skutecznie równoważ nierównomierne obciążenia.
Debugowanie i równoważenie obciążenia kodu równoległego to bezpłatna lekcja R Academy na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej R Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs R Academy zawiera 4 lekcji w sumie.
Dlaczego debugowanie równoległe jest trudne
Debugowanie kodu równoległego jest trudne, ponieważ procesy robocze działają w oddzielnych procesach: instrukcje print() nie są widoczne w procesie głównym, interaktywne debugery, takie jak browser(), nie działają wewnątrz procesów roboczych, a błędy są szeregowane i zgłaszane ponownie w procesie głównym, przez co tracą swój pierwotny ślad stosu.
library(parallel)
# Problem: cat() inside worker is invisible in the master
cl <- makeCluster(2)
clusterExport(cl, c())
result <- parLapply(cl, 1:4, function(i) {
# This cat() goes to the worker's stdout, NOT the master
cat('Worker processing', i, '\n') # you will NOT see this
i^2
})
cat('Master received:', unlist(result), '\n')
stopCluster(cl)Propagowanie błędów w parLapply
Gdy proces roboczy zgłosi błąd, parLapply() opakowuje go i zgłasza ponownie w procesie głównym. Komunikat błędu zostaje zachowany, ale zdalny ślad stosu nie. Przed zrównolegleniem należy zawsze najpierw przetestować funkcję sekwencyjnie.
library(parallel)
cl <- makeCluster(2)
# Step 1: test sequentially first!
my_fn <- function(x) {
if (x == 3) stop('bad value: 3')
sqrt(x)
}
# Sequential test reveals the bug before parallelising
tryCatch(
lapply(1:5, my_fn),
error = function(e) cat('Sequential error caught:', e$message, '\n')
)
# Fix: guard the bad case
my_fn_safe <- function(x) {
if (x <= 0) return(NA_real_)
sqrt(x)
}
clusterExport(cl, 'my_fn_safe')
cat(unlist(parLapply(cl, 1:5, my_fn_safe)), '\n')
stopCluster(cl)tryCatch wewnątrz funkcji procesów roboczych
Opakowanie treści funkcji procesu roboczego w tryCatch() pozwala przechwytywać i rejestrować błędy dla poszczególnych elementów bez przerywania całego zadania równoległego. W przypadku niepowodzenia należy zwrócić wartość sygnalizującą, na przykład NA, aby etap przetwarzania wyników mógł zidentyfikować problematyczne dane wejściowe.
library(parallel)
cl <- makeCluster(2)
safe_compute <- function(x) {
tryCatch({
if (x == 3) stop('simulated error on element 3')
list(value = x^2, error = NULL)
}, error = function(e) {
list(value = NA_real_, error = conditionMessage(e))
})
}
clusterExport(cl, 'safe_compute')
results <- parLapply(cl, 1:5, safe_compute)
for (i in seq_along(results)) {
r <- results[[i]]
if (is.null(r$error)) {
cat('Element', i, ': value =', r$value, '\n')
} else {
cat('Element', i, ': ERROR -', r$error, '\n')
}
}
stopCluster(cl)future::value() z tryCatch
W pakiecie future funkcja value(f) zgłasza ponownie zdalny błąd. Opakowanie jej w tryCatch() pozwala obsługiwać pojedyncze niepowodzenia obiektów future i jednocześnie zbierać wyniki z pozostałych obiektów future.
library(future)
plan(multisession, workers = 2)
# Create futures — some will fail
inputs <- list(4, -1, 9, 'x', 16)
futures <- lapply(inputs, function(val) {
future({
if (!is.numeric(val)) stop('non-numeric input')
if (val < 0) stop('negative input')
sqrt(val)
})
})
# Collect with individual error handling
results <- lapply(futures, function(f) {
tryCatch(
value(f),
error = function(e) paste('ERROR:', e$message)
)
})
for (i in seq_along(results)) {
cat('Input:', inputs[[i]], '-> Result:', as.character(results[[i]]), '\n')
}
plan(sequential)foreach z %dopar% i .errorhandling
Operator %dopar% z pakietu foreach rozdziela iteracje między zarejestrowany backend. Argument .errorhandling określa zachowanie w przypadku błędów: 'stop' (domyślnie — zatrzymanie), 'remove' (pominięcie) lub 'pass' (dołączenie obiektu warunku).
library(foreach)
library(doParallel)
cl <- makeCluster(2)
registerDoParallel(cl)
# .errorhandling = 'pass': failed elements return the condition
results <- foreach(
x = 1:6,
.errorhandling = 'pass'
) %dopar% {
if (x == 4) stop('bad element')
x^2
}
for (i in seq_along(results)) {
if (inherits(results[[i]], 'error')) {
cat('Element', i, ': ERROR -', results[[i]]$message, '\n')
} else {
cat('Element', i, ': value =', results[[i]], '\n')
}
}
stopCluster(cl)Równoważenie obciążenia: statyczne a dynamiczne
Planowanie statyczne z góry przydziela procesom roboczym fragmenty o takiej samej wielkości. Planowanie dynamiczne przydziela każdemu procesowi roboczemu jedno zadanie naraz, dzięki czemu szybkie procesy mogą pobierać kolejne zadania. Planowanie dynamiczne sprawdza się lepiej, gdy czasy wykonywania zadań znacznie się różnią.
library(parallel)
cl <- makeCluster(2)
# Simulate unequal task durations (element i takes i*0.05 seconds)
unequal_task <- function(i) {
Sys.sleep(i * 0.05)
i
}
clusterExport(cl, 'unequal_task')
# Static: parLapply distributes in fixed chunks
static_time <- system.time(
parLapply(cl, 1:6, unequal_task)
)[['elapsed']]
# Dynamic: clusterApplyLB assigns one task per available worker
dynamic_time <- system.time(
clusterApplyLB(cl, 1:6, unequal_task)
)[['elapsed']]
cat('Static LB: ', round(static_time, 2), 's\n')
cat('Dynamic LB:', round(dynamic_time, 2), 's\n')
stopCluster(cl)Dzielenie na fragmenty: zmniejszanie narzutu
Komunikacja między procesami wiąże się ze stałym narzutem dla każdego zadania. Podczas przetwarzania wielu małych elementów należy grupować je w większe fragmenty, aby każde wywołanie procesu roboczego wykonywało więcej pracy w ramach jednego komunikatu, zmniejszając stosunek narzutu do czasu obliczeń.
library(parallel)
cl <- makeCluster(2)
# Naive: 1000 tiny tasks — high overhead
tiny_task <- function(x) x^2
clusterExport(cl, 'tiny_task')
t1 <- system.time(parLapply(cl, 1:1000, tiny_task))[['elapsed']]
# Chunked: 10 tasks of 100 items each — low overhead
chunk_task <- function(chunk) sapply(chunk, function(x) x^2)
chunks <- split(1:1000, cut(1:1000, 10, labels = FALSE))
clusterExport(cl, 'chunk_task')
t2 <- system.time(parLapply(cl, chunks, chunk_task))[['elapsed']]
cat('Unchunked:', round(t1, 4), 's\n')
cat('Chunked: ', round(t2, 4), 's\n')
stopCluster(cl)Unikanie przesyłania dużych obiektów
Szeregowanie dużych obiektów (ramek danych, macierzy, modeli) do procesów roboczych jest kosztowne. Zamiast tego należy przekazywać tylko indeksy i odczytywać dane ze współdzielonego źródła (pliku, bazy danych lub danych wstępnie załadowanych w procesie roboczym). Proszę utrzymywać mały rozmiar danych przesyłanych do procesów roboczych.
library(parallel)
cl <- makeCluster(2)
# BAD: sends the entire large data frame to every worker
big_df <- data.frame(x = rnorm(1e5), y = rnorm(1e5))
clusterExport(cl, 'big_df') # 800 KB shipped to each worker
# BETTER: workers generate/read their own data slice
clusterEvalQ(cl, {
set.seed(Sys.getpid()) # unique per worker
local_data <- data.frame(x = rnorm(500), y = rnorm(500))
})
# Each worker uses its own local_data without receiving it from master
result <- parLapply(cl, 1:2, function(i) {
coef(lm(y ~ x, data = local_data))
})
print(result)
stopCluster(cl)Rejestrowanie komunikatów z procesów roboczych
Ponieważ procesy robocze nie mogą zapisywać w konsoli procesu głównego, należy przekierować ich dane wyjściowe do osobnych plików dziennika za pomocą makeCluster(outfile = '/path/to/log'). W systemie Unix można przekierować je do /dev/null lub do osobnego pliku dziennika dla każdego procesu roboczego.
library(parallel)
# Route all worker stdout/stderr to a log file
log_file <- tempfile(fileext = '.log')
cl <- makeCluster(2, outfile = log_file)
clusterEvalQ(cl, {
cat('[Worker', Sys.getpid(), '] started\n')
})
parLapply(cl, 1:4, function(i) {
cat('[Worker] processing item', i, '\n') # goes to log file
i * 10
})
stopCluster(cl)
# Read the log
log_content <- readLines(log_file)
cat(head(log_content, 8), sep = '\n')Profilowanie kodu równoległego
Proszę użyć system.time() do zmierzenia całkowitego czasu rzeczywistego i ręcznie rozdzielić narzut procesów roboczych od czasu obliczeń. Aby uzyskać więcej szczegółów, należy najpierw uruchomić docelową funkcję sekwencyjnie za pomocą profvis::profvis(), a dopiero potem zrównoleglić wąskie gardło.
library(parallel)
# Profile the task sequentially first
target_fn <- function(n) {
x <- rnorm(n)
list(
mean = mean(x),
sd = sd(x),
q95 = quantile(x, 0.95)
)
}
# Measure sequential baseline
t_seq <- system.time(lapply(rep(1000, 20), target_fn))[['elapsed']]
# Measure parallel speedup
cl <- makeCluster(2)
clusterExport(cl, 'target_fn')
t_par <- system.time(parLapply(cl, rep(1000, 20), target_fn))[['elapsed']]
stopCluster(cl)
cat('Sequential:', round(t_seq, 4), 's\n')
cat('Parallel: ', round(t_par, 4), 's\n')
cat('Efficiency:', round(t_seq / (t_par * 2) * 100, 1), '%\n')Kompletny odporny wzorzec
Połączenie wszystkich najlepszych praktyk: podział pracy na fragmenty, użycie równoważenia obciążenia, obsługa błędów dla poszczególnych elementów, rejestrowanie danych do pliku oraz zagwarantowanie zamknięcia klastra za pomocą on.exit().
library(parallel)
robust_parallel <- function(items, fn, workers = 2) {
cl <- makeCluster(workers, outfile = tempfile())
on.exit(stopCluster(cl), add = TRUE)
clusterExport(cl, 'fn', envir = environment())
safe_fn <- function(x) {
tryCatch(
fn(x),
error = function(e) list(error = e$message, value = NA)
)
}
clusterExport(cl, 'safe_fn', envir = environment())
# Load-balanced assignment for unequal task times
clusterApplyLB(cl, items, safe_fn)
}
results <- robust_parallel(1:8, function(x) {
if (x == 5) stop('deliberate error')
x^3
})
cat('Results:\n')
for (i in seq_along(results)) {
cat(' [', i, ']', ifelse(is.na(results[[i]]$value),
paste('ERROR:', results[[i]]$error),
results[[i]]
), '\n')
}Szybkie sprawdzenie
Mają Państwo 100 zadań równoległych o bardzo zróżnicowanych czasach wykonywania (niektóre trwają 10 ms, a inne 500 ms). Która strategia planowania zminimalizuje całkowity czas rzeczywisty?
Podsumowanie: debugowanie kodu równoległego
Najważniejsze informacje:
- Przed zrównolegleniem należy testować funkcje sekwencyjnie — najpierw użyć
lapply() - Proszę opakować treść funkcji procesów roboczych w
tryCatch(), aby przechwytywać błędy poszczególnych elementów bez przerywania zadania future::value()wraz ztryCatch()umożliwia sprawną obsługę pojedynczych niepowodzeń obiektów futureforeach %dopar%z.errorhandling = 'pass'zwraca obiekty błędów na liście wyników- Proszę używać
clusterApplyLB()do dynamicznego równoważenia obciążenia przy nierównych czasach wykonywania zadań - Należy dzielić małe zadania na fragmenty, aby zmniejszyć narzut komunikacji
- Rozmiar danych przesyłanych do procesów roboczych powinien być minimalny — należy przekazywać indeksy, a nie duże obiekty danych
- Dane wyjściowe procesów roboczych należy rejestrować w plikach za pomocą
makeCluster(outfile=)
# Canonical debugging workflow
# 1. Test sequentially
# lapply(inputs, my_fn)
# 2. Wrap in tryCatch
# safe_fn <- function(x) tryCatch(my_fn(x), error = function(e) NA)
# 3. Parallelise the safe version
# cl <- makeCluster(2); on.exit(stopCluster(cl))
# parLapply(cl, inputs, safe_fn)
cat('Debugging workflow: sequential -> tryCatch -> parallel\n')Często zadawane pytania
Czy lekcja „Debugowanie i równoważenie obciążenia kodu równoległego” jest bezpłatna?
Tak — pełny tekst „Debugowanie i równoważenie obciążenia kodu równoległego” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu R Academy, przejdź na CoddyKit PRO. Kurs R Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Debugowanie i równoważenie obciążenia kodu równoległego”?
Obsługuj błędy w procesach roboczych i skutecznie równoważ nierównomierne obciążenia. Ćwiczysz R Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć R Academy?
Nie wymagamy żadnego doświadczenia. R Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.
Ile czasu zajmuje lekcja „Debugowanie i równoważenie obciążenia kodu równoległego”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji R Academy?
Tak. Każda lekcja R Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Pakiet parallel i detectCores()
- Framework future
- furrr: równoległe operacje purrr
- Debugowanie i równoważenie obciążenia kodu równoległego