0Pricing
R Academy · Урок

Отладка и балансировка нагрузки параллельного кода

Обрабатывайте ошибки в параллельных рабочих процессах и эффективно распределяйте неравномерную нагрузку

«Отладка и балансировка нагрузки параллельного кода» — бесплатный урок R Academy на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения R Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс R Academy содержит 4 уроков всего.

Почему отладка параллельного кода сложна

Отлаживать параллельный код сложно, поскольку рабочие процессы выполняются в отдельных процессах: инструкции print() невидимы главному процессу, интерактивные отладчики вроде browser() не работают внутри рабочих процессов, а ошибки сериализуются и повторно вызываются в главном процессе, из-за чего теряется исходная трассировка стека.

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)

Передача ошибок в parLapply

Если рабочий процесс вызывает ошибку, parLapply() оборачивает её и повторно вызывает в главном процессе. Сообщение об ошибке сохраняется, но удалённая трассировка стека — нет. Всегда сначала проверяйте функцию последовательно, а уже затем распараллеливайте её.

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 внутри функций рабочих процессов

Обёртывание тела функции рабочего процесса в tryCatch() позволяет перехватывать ошибки и вести журнал для каждого элемента, не прерывая всю параллельную задачу. При сбое возвращайте специальное значение, например NA, чтобы на этапе последующей обработки можно было выявить проблемные входные данные.

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

В пакете future функция value(f) повторно вызывает удалённую ошибку. Обёртывание в tryCatch() позволяет обрабатывать отдельные сбои future и продолжать сбор результатов остальных 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 с %dopar% и .errorhandling

Оператор %dopar% пакета foreach распределяет итерации между зарегистрированными обработчиками. Аргумент .errorhandling определяет поведение при ошибках: 'stop' (по умолчанию, остановить), 'remove' (пропустить) или 'pass' (включить объект условия).

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)

Балансировка нагрузки: статическая и динамическая

Статическое планирование заранее распределяет между рабочими процессами фрагменты одинакового размера. Динамическое планирование выдаёт каждому рабочему процессу по одной задаче, поэтому быстрые процессы получают дополнительные задачи. Динамический вариант лучше, когда длительность задач сильно различается.

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)

Разбиение на фрагменты: уменьшение накладных расходов

Межпроцессное взаимодействие имеет фиксированные накладные расходы для каждой задачи. При обработке множества небольших элементов объединяйте их в более крупные фрагменты, чтобы каждый вызов рабочего процесса выполнял больше работы за одно сообщение и уменьшал отношение накладных расходов к вычислениям.

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)

Не отправляйте большие объекты

Сериализация больших объектов (таблиц данных, матриц, моделей) для рабочих процессов требует значительных ресурсов. Вместо этого передавайте только индексы и считывайте данные из общего источника (файла, базы данных или данных, предварительно загруженных в рабочий процесс). Старайтесь, чтобы передаваемые рабочим процессам данные были небольшими.

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)

Ведение журнала в рабочих процессах

Поскольку рабочие процессы не могут записывать данные в консоль главного процесса, перенаправляйте их вывод в отдельные файлы журнала с помощью makeCluster(outfile = '/path/to/log'). В Unix вывод можно перенаправить в /dev/null или в отдельный файл журнала для каждого рабочего процесса.

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

Профилирование параллельного кода

Используйте system.time(), чтобы измерить общее реальное время выполнения, а затем вручную разделите накладные расходы рабочих процессов и время вычислений. Для более подробного анализа сначала запустите целевую функцию последовательно с помощью profvis::profvis(), а затем распараллельте узкое место.

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

Полный надёжный шаблон

Объедините все лучшие практики: разбейте работу на фрагменты, используйте сбалансированное распределение нагрузки, обрабатывайте ошибки для каждого элемента, записывайте сведения в файл и гарантируйте очистку кластера с помощью 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')
}

Быстрая проверка

У вас есть 100 параллельных задач с сильно различающимся временем выполнения: некоторые занимают 10 мс, а другие — 500 мс. Какая стратегия планирования минимизирует общее реальное время выполнения?

Итоги: отладка параллельного кода

Главное:

  • Проверяйте функции последовательно до распараллеливания — сначала используйте lapply()
  • Оборачивайте тела функций рабочих процессов в tryCatch(), чтобы перехватывать ошибки отдельных элементов и не прерывать задачу
  • future::value() вместе с tryCatch() позволяет корректно обрабатывать отдельные сбои future
  • foreach %dopar% с .errorhandling = 'pass' возвращает объекты ошибок в списке результатов
  • Используйте clusterApplyLB() для динамической балансировки нагрузки при разной длительности задач
  • Объединяйте небольшие задачи во фрагменты, чтобы уменьшить накладные расходы на обмен данными
  • Сводите передаваемые рабочим процессам данные к минимуму: передавайте индексы, а не большие объекты данных
  • Записывайте вывод рабочих процессов в файлы с помощью 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')

Часто задаваемые вопросы

Урок «Отладка и балансировка нагрузки параллельного кода» бесплатный?

Да — полный текст урока «Отладка и балансировка нагрузки параллельного кода» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс R Academy, подпишись на CoddyKit PRO. Курс R Academy содержит 4 уроков всего.

Чему я научусь в уроке «Отладка и балансировка нагрузки параллельного кода»?

Обрабатывайте ошибки в параллельных рабочих процессах и эффективно распределяйте неравномерную нагрузку Ты практикуешь R Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать R Academy?

Предыдущий опыт не требуется. R Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.

Сколько времени занимает урок «Отладка и балансировка нагрузки параллельного кода»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке R Academy?

Да. Каждый урок R Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Пакет parallel и detectCores()
  2. Фреймворк future
  3. furrr: параллельные операции purrr
  4. Отладка и балансировка нагрузки параллельного кода
← Назад к R Academy