0Pricing
R Academy · Lezione

Debugging e bilanciamento del carico nel codice parallelo

Gestisca gli errori nei worker paralleli e bilanci efficacemente carichi di lavoro non uniformi

Debugging e bilanciamento del carico nel codice parallelo è una lezione R Academy gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento R Academy, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso R Academy include 4 lezioni in totale.

Perché il debugging parallelo è difficile

Il debugging del codice parallelo è complesso perché i worker vengono eseguiti in processi separati: le istruzioni print() non sono visibili al master, i debugger interattivi come browser() non funzionano all'interno dei worker e gli errori vengono serializzati e generati nuovamente nel master, perdendo la traccia dello stack originale.

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)

Propagazione degli errori in parLapply

Quando un worker genera un errore, parLapply() lo racchiude e lo genera nuovamente nel master. Il messaggio di errore viene conservato, ma la traccia dello stack remoto no. Verifichi sempre prima la funzione in sequenza, prima di parallelizzarla.

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 nelle funzioni worker

Racchiudere il corpo della funzione worker in tryCatch() consente di acquisire e registrare gli errori per ogni elemento senza interrompere l'intero processo parallelo. In caso di errore, restituisca un valore sentinella, ad esempio NA, in modo che la post-elaborazione possa individuare gli input problematici.

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() con tryCatch

Nel pacchetto future, value(f) genera nuovamente l'errore remoto. Racchiuderlo in tryCatch() consente di gestire i singoli errori dei future continuando a raccogliere i risultati degli altri 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 con %dopar% e .errorhandling

L'operatore %dopar% del pacchetto foreach distribuisce le iterazioni tra i worker di un backend registrato. L'argomento .errorhandling controlla cosa accade in caso di errore: 'stop' (predefinito), 'remove' (ignora) o 'pass' (include l'oggetto condizione).

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)

Bilanciamento del carico: statico e dinamico

La pianificazione statica assegna in anticipo ai worker blocchi di dimensioni uguali. La pianificazione dinamica assegna a ogni worker un'attività alla volta, permettendo ai worker più veloci di prenderne altre. La modalità dinamica è migliore quando la durata delle attività varia molto.

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)

Suddivisione in blocchi: ridurre l'overhead

La comunicazione tra processi comporta un overhead fisso per ogni attività. Quando si elaborano molti elementi piccoli, li si raggruppi in blocchi più grandi, così ogni chiamata al worker esegue più lavoro per messaggio, riducendo il rapporto tra overhead e calcolo.

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)

Evitare l'invio di oggetti di grandi dimensioni

La serializzazione di oggetti di grandi dimensioni, come data frame, matrici e modelli, verso i worker è costosa. Passi invece solo gli indici e legga i dati da una fonte condivisa, come un file o un database, oppure da dati precaricati nel worker. Mantenga ridotti i dati inviati ai worker.

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)

Registrazione dai worker

Poiché i worker non possono scrivere nella console del master, reindirizzi l'output dei worker verso file di log separati usando makeCluster(outfile = '/path/to/log'). Su Unix è possibile reindirizzarlo a /dev/null o a un file di log dedicato per ogni worker.

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')

Profilazione del codice parallelo

Usi system.time() per misurare il tempo totale trascorso e separi manualmente l'overhead dei worker dal tempo di calcolo. Per maggiori dettagli, esegua prima la funzione interessata in sequenza con profvis::profvis(), quindi parallelizzi il collo di bottiglia.

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')

Schema completo e robusto

Combinando tutte le pratiche consigliate: suddividere il lavoro in blocchi, usare un'assegnazione con bilanciamento del carico, gestire gli errori per elemento, registrare l'output su file e garantire la pulizia del cluster con 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')
}

Verifica rapida

Dispone di 100 attività parallele con tempi di esecuzione molto variabili (alcune durano 10 ms, altre 500 ms). Quale strategia di pianificazione ridurrà al minimo il tempo totale trascorso?

Riepilogo: debugging del codice parallelo

Punti chiave:

  • Verifichi le funzioni in sequenza prima di parallelizzarle: usi prima lapply()
  • Racchiuda i corpi dei worker in tryCatch() per acquisire gli errori dei singoli elementi senza interrompere il processo
  • future::value() con tryCatch() gestisce in modo sicuro i singoli errori dei future
  • foreach %dopar% con .errorhandling = 'pass' restituisce gli oggetti errore nella lista dei risultati
  • Usi clusterApplyLB() per il bilanciamento dinamico del carico con attività di durata diversa
  • Raggruppi le attività piccole in blocchi per ridurre l'overhead di comunicazione
  • Mantenga minimi i dati inviati ai worker: passi indici, non oggetti di dati voluminosi
  • Registri l'output dei worker su file tramite 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')

Domande Frequenti

La lezione «Debugging e bilanciamento del carico nel codice parallelo» è gratuita?

Sì — il testo completo di «Debugging e bilanciamento del carico nel codice parallelo» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso R Academy, passa a CoddyKit PRO. Il corso R Academy include 4 lezioni in totale.

Cosa imparerò in «Debugging e bilanciamento del carico nel codice parallelo»?

Gestisca gli errori nei worker paralleli e bilanci efficacemente carichi di lavoro non uniformi Eserciti R Academy con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare R Academy?

Non è richiesta alcuna esperienza precedente. R Academy su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «Debugging e bilanciamento del carico nel codice parallelo»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione R Academy?

Sì. Ogni lezione R Academy include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Pacchetto parallel e detectCores()
  2. Il framework future
  3. furrr: operazioni parallele di purrr
  4. Debugging e bilanciamento del carico nel codice parallelo
← Torna a R Academy