0Pricing
R Academy · Leçon

Déboguer et équilibrer la charge du code parallèle

Gérez les erreurs des travailleurs parallèles et répartissez efficacement les charges de travail inégales.

Déboguer et équilibrer la charge du code parallèle est une leçon R Academy gratuite sur CoddyKit. Ceci est la leçon 4 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage R Academy, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours R Academy comprend 4 leçons au total.

Pourquoi le débogage parallèle est difficile

Le débogage du code parallèle est difficile, car les processus de travail s’exécutent dans des processus distincts : les instructions print() sont invisibles dans la console principale, les débogueurs interactifs comme browser() ne fonctionnent pas dans les processus de travail, et les erreurs sont sérialisées puis relancées dans le processus principal, ce qui fait perdre leur trace d’appels d’origine.

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)

Propagation des erreurs dans parLapply

Lorsqu’un processus de travail génère une erreur, parLapply() l’enveloppe et la relance dans le processus principal. Le message d’erreur est conservé, mais pas la trace d’appels distante. Testez toujours votre fonction séquentiellement avant de la paralléliser.

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 dans les fonctions de travail

Entourer le corps de la fonction de travail avec tryCatch() permet de capturer et d’enregistrer les erreurs pour chaque élément sans interrompre toute la tâche parallèle. En cas d’échec, renvoyez une valeur indicatrice (par exemple NA) afin que le post-traitement puisse identifier les entrées problématiques.

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

Dans le package future, value(f) relance l’erreur distante. L’entourer de tryCatch() permet de gérer les échecs individuels des futures tout en continuant à recueillir les résultats des autres futures.

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 avec %dopar% et .errorhandling

L’opérateur %dopar% du package foreach répartit les itérations entre les processus d’un moteur enregistré. L’argument .errorhandling contrôle le comportement en cas d’erreur : 'stop' (par défaut), 'remove' (ignorer) ou 'pass' (inclure l’objet condition).

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épartition de la charge : statique ou dynamique

La planification statique attribue à l’avance aux processus de travail des blocs de taille égale. La planification dynamique donne une tâche à la fois à chaque processus ; les processus rapides en prennent donc davantage. La méthode dynamique est préférable lorsque la durée des tâches varie beaucoup.

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)

Regroupement en blocs : réduire les coûts supplémentaires

La communication entre processus entraîne un coût fixe par tâche. Lorsque vous traitez de nombreux petits éléments, regroupez-les en blocs plus importants afin que chaque appel à un processus de travail effectue davantage de travail par message, ce qui réduit le rapport entre les coûts de communication et le calcul.

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)

Éviter d’envoyer de gros objets

La sérialisation de gros objets (trames de données, matrices, modèles) vers les processus de travail est coûteuse. Transmettez plutôt uniquement les indices et lisez les données depuis une source partagée (fichier, base de données ou données préchargées dans le processus de travail). Gardez les données transmises aux processus de travail peu volumineuses.

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)

Journalisation depuis les processus de travail

Comme les processus de travail ne peuvent pas écrire dans la console principale, redirigez leur sortie vers des fichiers journaux propres à chaque processus avec makeCluster(outfile = '/path/to/log'). Sous Unix, vous pouvez rediriger la sortie vers /dev/null ou vers un fichier journal dédié à chaque processus.

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

Analyser les performances du code parallèle

Utilisez system.time() pour mesurer le temps réel total, puis décomposez manuellement le coût des processus de travail et le temps de calcul. Pour plus de détails, exécutez d’abord séquentiellement la fonction cible avec profvis::profvis(), puis parallélisez le goulet d’étranglement.

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

Modèle robuste complet

En réunissant toutes les bonnes pratiques : regrouper le travail en blocs, utiliser une répartition équilibrée de la charge, gérer les erreurs élément par élément, écrire les journaux dans un fichier et garantir le nettoyage du groupe de processus avec 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')
}

Vérification rapide

Vous avez 100 tâches parallèles dont les durées d’exécution sont très variables (certaines durent 10 ms, d’autres 500 ms). Quelle stratégie de planification minimisera le temps réel total ?

Récapitulatif : déboguer du code parallèle

À retenir :

  • Testez les fonctions séquentiellement avant de les paralléliser : utilisez d’abord lapply()
  • Entourez le corps des fonctions de travail de tryCatch() pour capturer les erreurs par élément sans interrompre la tâche
  • future::value() avec tryCatch() gère correctement les échecs individuels des futures
  • foreach %dopar% avec .errorhandling = 'pass' renvoie les objets d’erreur dans la liste de résultats
  • Utilisez clusterApplyLB() pour une répartition dynamique de la charge lorsque les durées des tâches sont inégales
  • Regroupez les petites tâches en blocs pour réduire les coûts de communication
  • Limitez les données transmises aux processus de travail : transmettez des indices, pas de gros objets de données
  • Écrivez la sortie des processus de travail dans des fichiers avec 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')

Questions Fréquemment Posées

La leçon « Déboguer et équilibrer la charge du code parallèle » est-elle gratuite ?

Oui — le texte complet de « Déboguer et équilibrer la charge du code parallèle » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours R Academy, passe à CoddyKit PRO. Le cours R Academy comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « Déboguer et équilibrer la charge du code parallèle » ?

Gérez les erreurs des travailleurs parallèles et répartissez efficacement les charges de travail inégales. Tu pratiques R Academy avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer R Academy ?

Aucune expérience préalable n'est requise. R Academy sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 4 sur 4.

Combien de temps prend la leçon « Déboguer et équilibrer la charge du code parallèle » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon R Academy ?

Oui. Chaque leçon R Academy inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. Paquet parallel et detectCores()
  2. Le cadriciel future
  3. furrr : opérations parallèles avec purrr
  4. Déboguer et équilibrer la charge du code parallèle
← Retour à R Academy