Depuração e balanceamento de carga de código paralelo
Lide com erros nos trabalhadores paralelos e equilibre cargas de trabalho desiguais com eficiência.
Depuração e balanceamento de carga de código paralelo é uma aula grátis de R Academy no CoddyKit. Esta é a aula 4 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de R Academy, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de R Academy inclui 4 aulas no total.
Por que a depuração paralela é difícil
Depurar código paralelo é desafiador porque os trabalhadores são executados em processos separados: as instruções print() não ficam visíveis para o processo principal, depuradores interativos como browser() não funcionam dentro dos trabalhadores e os erros são serializados e lançados novamente no processo principal, perdendo seu rastreamento de pilha 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)Propagação de erros em parLapply
Quando um trabalhador gera um erro, parLapply() o encapsula e o lança novamente no processo principal. A mensagem de erro é preservada, mas o rastreamento de pilha remoto não. Sempre teste sua função sequencialmente primeiro, antes de paralelizá-la.
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 das funções dos trabalhadores
Envolver o corpo da função do trabalhador em tryCatch() permite capturar e registrar erros por elemento sem interromper toda a tarefa paralela. Em caso de falha, retorne um valor sentinela (por exemplo, NA) para que o pós-processamento possa identificar as 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() com tryCatch
No pacote future, value(f) lança novamente o erro remoto. Envolvê-lo em tryCatch() permite tratar falhas individuais de futures enquanto o senhor continua coletando os resultados dos outros 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 com %dopar% e .errorhandling
O operador %dopar% do pacote foreach distribui as iterações entre um mecanismo registrado. O argumento .errorhandling controla o que acontece em caso de erros: 'stop' (padrão), 'remove' (ignorar) ou 'pass' (incluir o objeto de condição).
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)Balanceamento de carga: estático versus dinâmico
O agendamento estático pré-atribui blocos de tamanho igual aos trabalhadores. O agendamento dinâmico fornece uma tarefa por vez a cada trabalhador, permitindo que os trabalhadores rápidos assumam mais tarefas. O modo dinâmico é melhor quando a duração das tarefas varia muito.
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)Agrupamento em blocos: reduzindo a sobrecarga
A comunicação entre processos tem uma sobrecarga fixa por tarefa. Ao processar muitos itens pequenos, agrupe-os em blocos maiores para que cada chamada do trabalhador faça mais trabalho por mensagem, reduzindo a proporção entre sobrecarga e 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 (quadros de dados, matrizes, modelos) para os trabalhadores é dispendioso. Em vez disso, passe apenas os índices e leia os dados de uma fonte compartilhada (arquivo, banco de dados ou dados pré-carregados no trabalhador). Mantenha pequenas as cargas de dados dos trabalhadores.
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 de atividades dos trabalhadores
Como os trabalhadores não podem escrever no console do processo principal, redirecione a saída de cada trabalhador para arquivos de registro individuais usando makeCluster(outfile = '/path/to/log'). No Unix, o senhor pode redirecionar a saída para /dev/null ou para um arquivo de registro específico de cada trabalhador.
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')Criação de perfil de código paralelo
Use system.time() para medir o tempo total de execução, e decomponha manualmente o tempo de sobrecarga dos trabalhadores e o tempo de cálculo. Para obter mais detalhes, execute primeiro a função-alvo sequencialmente com profvis::profvis() e só depois paralelize o gargalo.
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')Padrão robusto completo
Combinando todas as práticas recomendadas: divida o trabalho em blocos, use atribuição com balanceamento de carga, trate os erros por elemento, registre as informações em arquivo e garanta a limpeza do cluster com 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ção rápida
O senhor tem 100 tarefas paralelas com tempos de execução muito variáveis (algumas levam 10 ms, outras 500 ms). Qual estratégia de agendamento minimizará o tempo total de execução?
Recapitulação: depuração de código paralelo
Principais conclusões:
- Teste as funções sequencialmente antes de paralelizá-las — use
lapply()primeiro - Envolva os corpos das funções dos trabalhadores em
tryCatch()para capturar erros por elemento sem interromper a tarefa future::value()comtryCatch()trata adequadamente falhas individuais de futuresforeach %dopar%com.errorhandling = 'pass'retorna objetos de erro na lista de resultados- Use
clusterApplyLB()para balanceamento de carga dinâmico com tarefas de durações diferentes - Agrupe tarefas pequenas em blocos para reduzir a sobrecarga de comunicação
- Mantenha mínimas as cargas de dados dos trabalhadores — passe índices, não objetos de dados grandes
- Registre a saída dos trabalhadores em arquivos por meio de
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')Perguntas Frequentes
A aula “Depuração e balanceamento de carga de código paralelo” é grátis?
Sim — o texto completo de “Depuração e balanceamento de carga de código paralelo” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de R Academy, atualize para CoddyKit PRO. O curso de R Academy inclui 4 aulas no total.
O que vou aprender em “Depuração e balanceamento de carga de código paralelo”?
Lide com erros nos trabalhadores paralelos e equilibre cargas de trabalho desiguais com eficiência. Você pratica R Academy com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.
Preciso ter experiência prévia para começar R Academy?
Nenhuma experiência prévia é necessária. R Academy no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 4 de 4.
Quanto tempo leva a aula “Depuração e balanceamento de carga de código paralelo”?
A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.
Posso escrever e executar código nesta aula de R Academy?
Sim. Cada aula de R Academy inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.
Todas as aulas deste curso
- Pacote parallel e detectCores()
- O framework future
- furrr: operações paralelas com purrr
- Depuração e balanceamento de carga de código paralelo