Rinnakkaisen koodin virheenkorjaus ja kuormantasaus
Käsitelkää rinnakkaisten työntekijöiden virheitä ja tasatkaa epätasaiset työkuormat tehokkaasti.
Rinnakkaisen koodin virheenkorjaus ja kuormantasaus on ilmainen R Academy-oppitunti CoddyKitissä. Tämä on oppitunti 4/4. Voit lukea tästä oppimispolusta kokonaan mitkä tahansa 3 oppituntia ilmaiseksi — sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä käytännön harjoittelun sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. Oppitunti kuuluu R Academy-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. R Academy-kurssilla on yhteensä 4 oppituntia.
Miksi rinnakkaisen koodin virheenkorjaus on vaikeaa
Rinnakkaisen koodin virheenkorjaus on haastavaa, koska workerit suoritetaan erillisissä prosesseissa: print()-lausekkeet eivät näy pääprosessissa, vuorovaikutteiset virheenkorjaimet, kuten browser(), eivät toimi workereissa, ja virheet sarjallistetaan ja välitetään uudelleen pääprosessissa, jolloin alkuperäinen kutsupino menetetään.
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)Virheiden välittäminen parLapply-funktiossa
Kun worker aiheuttaa virheen, parLapply() käsittelee sen ja välittää sen uudelleen pääprosessissa. Virheilmoitus säilyy, mutta etäprosessin kutsupino ei. Testaa funktiosi aina ensin peräkkäin ennen sen suorittamista rinnakkain.
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 worker-funktioiden sisällä
Worker-funktion rungon kääriminen tryCatch()-rakenteeseen mahdollistaa virheiden keräämisen ja kirjaamisen alkioittain ilman koko rinnakkaisen työn kaatumista. Palauta virheen sattuessa merkkiluku, esimerkiksi NA, jotta ongelmalliset syötteet voidaan tunnistaa jälkikäsittelyssä.
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() ja tryCatch
future-paketissa value(f) välittää etäprosessin virheen uudelleen kutsujalle. Käärimällä sen tryCatch()-rakenteeseen voit käsitellä yksittäisten futurejen epäonnistumiset ja jatkaa muiden futurejen tulosten keräämistä.
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 ja %dopar% sekä .errorhandling
foreach-paketin %dopar%-operaattori jakaa iteraatiot rekisteröidyn taustajärjestelmän kesken. .errorhandling-argumentti määrittää virheiden käsittelyn: 'stop' (oletus), 'remove' (ohita) tai 'pass' (sisällytä condition-objekti).
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)Kuormantasaus: staattinen ja dynaaminen
Staattisessa ajoituksessa workereille määritetään etukäteen samankokoiset tehtäväerät. Dynaamisessa ajoituksessa kukin worker saa yhden tehtävän kerrallaan, joten nopeat workerit poimivat niitä lisää. Dynaaminen ajoitus on parempi, kun tehtävien kestot vaihtelevat suuresti.
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)Lohkoihin jakaminen: yleiskustannusten vähentäminen
Prosessien väliseen viestintään liittyy kiinteä kustannus tehtävää kohden. Kun käsittelet monia pieniä alkioita, ryhmittele ne suuremmiksi lohkoiksi, jotta jokainen worker-kutsu tekee enemmän työtä viestiä kohden. Näin laskennan osuus kokonaiskustannuksista kasvaa.
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)Vältä suurten objektien lähettämistä
Suurten objektien, kuten data framejen, matriisien ja mallien, sarjallistaminen workereille on kallista. Välitä sen sijaan vain indeksit ja lue data jaetusta lähteestä, kuten tiedostosta tai tietokannasta, tai workerissa valmiiksi ladatusta lähteestä. Pidä workerien hyötykuormat pieninä.
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)Kirjaaminen workereista
Koska workerit eivät voi kirjoittaa pääprosessin konsoliin, ohjaa workerien tulosteet omiin lokitiedostoihinsa käyttämällä makeCluster(outfile = '/path/to/log')-kutsua. Unixissa voit ohjata tulosteen kohteeseen /dev/null tai kullekin workerille omaan lokitiedostoon.
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')Rinnakkaisen koodin profilointi
Mittaa kokonaiskestoa system.time()-funktiolla ja erottele workerien yleiskustannukset laskenta-ajasta manuaalisesti. Saat tarkempia tietoja suorittamalla kohdefunktion ensin peräkkäin komennolla profvis::profvis() ja rinnakkaistamalla sitten pullonkaulan.
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')Täydellinen vikasietoinen malli
Yhdistä kaikki parhaat käytännöt: jaa työ lohkoihin, käytä kuormatasapainotettua jakoa, käsittele virheet alkioittain, kirjaa tapahtumat tiedostoon ja varmista klusterin siivous on.exit()-kutsulla.
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')
}Pikatarkistus
Sinulla on 100 rinnakkaista tehtävää, joiden suoritusajat vaihtelevat huomattavasti (jotkin kestävät 10 ms ja toiset 500 ms). Mikä ajoitusstrategia minimoi kokonaisajan?
Kertaus: rinnakkaisen koodin virheenkorjaus
Tärkeimmät asiat:
- Testaa funktiot peräkkäin ennen rinnakkaistamista – käytä ensin
lapply()-funktiota - Kääri workerien rungot
tryCatch()-rakenteeseen, jotta voit kerätä alkiokohtaiset virheet työn kaatumatta future::value()yhdessätryCatch()-rakenteen kanssa käsittelee yksittäisten futurejen epäonnistumiset hallitustiforeach %dopar%ja.errorhandling = 'pass'palauttavat virheobjektit tuloslistassa- Käytä
clusterApplyLB()-funktiota dynaamiseen kuormantasaukseen, kun tehtävien kestot vaihtelevat - Jaa pienet tehtävät lohkoihin viestinnän yleiskustannusten vähentämiseksi
- Pidä workereille lähetettävä hyötykuorma mahdollisimman pienenä – välitä indeksit, älä suuria dataobjekteja
- Kirjaa workerien tulosteet tiedostoihin komennolla
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')Opi R tekoälytuutorin avulla — ilmaiseksi
Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.
- Kurssit
- 43
- Oppitunnit
- 159
Usein kysytyt kysymykset
Onko oppitunti ”Rinnakkaisen koodin virheenkorjaus ja kuormantasaus” ilmainen?
Kyllä — voit lukea täällä verkossa kokonaan ilmaiseksi mitkä tahansa R Academy-oppimispolun 3 oppituntia, myös oppitunnin “Rinnakkaisen koodin virheenkorjaus ja kuormantasaus”. Sen jälkeen CoddyKit PRO avaa kaikki oppitunnit sekä interaktiiviset harjoitukset sisäänrakennetulla koodieditorilla ja ympäri vuorokauden toimivalla tekoälytuutorilla. R Academy-kurssilla on yhteensä 4 oppituntia.
Mitä opin oppitunnilla ”Rinnakkaisen koodin virheenkorjaus ja kuormantasaus”?
Käsitelkää rinnakkaisten työntekijöiden virheitä ja tasatkaa epätasaiset työkuormat tehokkaasti. Harjoittelet R Academy-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.
Tarvitsenko kokemusta aloittaakseni R Academy-opiskelun?
Aiempi kokemus ei ole tarpeen. CoddyKitin R Academy-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 4/4.
Kuinka kauan ”Rinnakkaisen koodin virheenkorjaus ja kuormantasaus”-oppitunnin suorittaminen kestää?
Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.
Voinko kirjoittaa ja suorittaa koodia tällä R Academy-oppitunnilla?
Kyllä. Jokainen R Academy-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.
Kaikki tämän kurssin oppitunnit
- parallel-paketti ja detectCores()
- future-kehys
- furrr: purrr-operaatioiden rinnakkaistaminen
- Rinnakkaisen koodin virheenkorjaus ja kuormantasaus