0Pricing
R Academy · Lección

Depuración y equilibrio de carga del código paralelo

Gestione errores en los workers paralelos y equilibre eficazmente cargas de trabajo desiguales.

Depuración y equilibrio de carga del código paralelo es una lección gratuita de R Academy en CoddyKit. Esta es la lección 4 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de R Academy, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de R Academy incluye 4 lecciones en total.

Por qué es difícil depurar código paralelo

Depurar código paralelo es complicado porque los workers se ejecutan en procesos independientes: las instrucciones print() no son visibles para el proceso principal, los depuradores interactivos como browser() no funcionan dentro de los workers y los errores se serializan y vuelven a lanzarse en el proceso principal, por lo que se pierde su traza de pila original.

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)

Propagación de errores en parLapply

Cuando un worker produce un error, parLapply() lo envuelve y vuelve a lanzarlo en el proceso principal. El mensaje de error se conserva, pero no la traza de pila remota. Pruebe siempre primero su función secuencialmente antes de paralelizarla.

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 dentro de funciones worker

Envolver el cuerpo de la función worker en tryCatch() permite capturar y registrar los errores de cada elemento sin bloquear todo el trabajo paralelo. En caso de error, devuelva un valor indicador, como NA, para que el posprocesamiento pueda identificar las entradas problemáticas.

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

En el paquete future, value(f) vuelve a lanzar el error remoto. Envolverlo en tryCatch() permite gestionar los fallos de futures individuales mientras continúa recopilando los resultados de los demás 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 con %dopar% y .errorhandling

El operador %dopar% del paquete foreach distribuye las iteraciones entre un backend registrado. El argumento .errorhandling controla qué ocurre cuando hay errores: 'stop' (predeterminado), 'remove' (omitir) o 'pass' (incluir el objeto de condición).

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)

Equilibrio de carga: estático frente a dinámico

La planificación estática asigna previamente a los workers bloques del mismo tamaño. La planificación dinámica proporciona una tarea cada vez a cada worker, de modo que los workers rápidos pueden tomar más tareas. La planificación dinámica es mejor cuando las duraciones de las tareas varían mucho.

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)

Agrupación en bloques: reducción de la sobrecarga

La comunicación entre procesos tiene una sobrecarga fija por tarea. Al procesar muchos elementos pequeños, agrúpelos en bloques más grandes para que cada llamada a un worker realice más trabajo por mensaje y reduzca la proporción entre sobrecarga y cálculo.

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)

Evite enviar objetos grandes

Serializar objetos grandes (data frames, matrices y modelos) para enviarlos a los workers resulta costoso. En su lugar, pase solo los índices y lea los datos desde una fuente compartida (un archivo o una base de datos, o datos precargados en el worker). Mantenga pequeñas las cargas útiles de los workers.

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)

Registro desde los workers

Dado que los workers no pueden escribir en la consola del proceso principal, redirija su salida a archivos de registro independientes mediante makeCluster(outfile = '/path/to/log'). En Unix puede redirigirla a /dev/null o a un archivo de registro específico para cada 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')

Perfilado de código paralelo

Use system.time() para medir el tiempo total de ejecución y descomponga manualmente la sobrecarga de los workers y el tiempo de cálculo. Para obtener más detalles, ejecute primero la función objetivo secuencialmente con profvis::profvis() y, después, paralelice el cuello de botella.

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

Patrón robusto completo

Combine todas las prácticas recomendadas: divida el trabajo en bloques, use una asignación equilibrada según la carga, gestione los errores de cada elemento, registre la información en un archivo y garantice la limpieza del clúster 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')
}

Comprobación rápida

Tiene 100 tareas paralelas con tiempos de ejecución muy variables (algunas tardan 10 ms y otras, 500 ms). ¿Qué estrategia de planificación minimizará el tiempo total de ejecución?

Resumen: depuración de código paralelo

Ideas clave:

  • Pruebe las funciones secuencialmente antes de paralelizarlas; use primero lapply()
  • Envuelva los cuerpos de los workers en tryCatch() para capturar los errores de cada elemento sin bloquear el trabajo
  • future::value() con tryCatch() gestiona correctamente los fallos de futures individuales
  • foreach %dopar% con .errorhandling = 'pass' devuelve objetos de error en la lista de resultados
  • Use clusterApplyLB() para equilibrar dinámicamente la carga cuando las duraciones de las tareas son desiguales
  • Agrupe las tareas pequeñas en bloques para reducir la sobrecarga de comunicación
  • Mantenga mínimas las cargas útiles de los workers: pase índices, no objetos de datos grandes
  • Registre la salida de los workers en archivos mediante 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')

Preguntas frecuentes

¿La lección «Depuración y equilibrio de carga del código paralelo» es gratis?

Sí — el texto completo de «Depuración y equilibrio de carga del código paralelo» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de R Academy, actualiza a CoddyKit PRO. El curso de R Academy incluye 4 lecciones en total.

¿Qué aprenderé en «Depuración y equilibrio de carga del código paralelo»?

Gestione errores en los workers paralelos y equilibre eficazmente cargas de trabajo desiguales. Practicas R Academy con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar R Academy?

No se requiere experiencia previa. R Academy en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 4 de 4.

¿Cuánto tiempo toma la lección «Depuración y equilibrio de carga del código paralelo»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de R Academy?

Sí. Cada lección de R Academy incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Paquete parallel y detectCores()
  2. El framework future
  3. furrr: operaciones paralelas de purrr
  4. Depuración y equilibrio de carga del código paralelo
← Volver a R Academy