0Pricing
R Academy · レッスン

並列コードのデバッグと負荷分散

並列ワーカーで発生するエラーを処理し、不均衡なワークロードを効果的に分散します。

「並列コードのデバッグと負荷分散」はCoddyKit上の無料R Academyレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これは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)

tryCatchを使ったfuture::value()

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

foreachパッケージの%dopar%演算子は、登録済みのバックエンドに反復処理を分散します。.errorhandling引数は、エラー発生時の動作を制御します。'stop'(デフォルト)、'remove'(スキップ)、'pass'(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)

負荷分散: 静的と動的

静的スケジューリングでは、同じサイズのチャンクをあらかじめワーカーに割り当てます。動的スケジューリングでは、各ワーカーに一度に1つずつタスクを与えるため、処理の速いワーカーがより多くのタスクを引き受けます。タスクの実行時間に大きなばらつきがある場合は、動的スケジューリングが適しています。

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)

チャンク化: オーバーヘッドの削減

プロセス間通信には、タスクごとに固定のオーバーヘッドがあります。小さな項目を大量に処理する場合は、より大きなチャンクにまとめます。そうすると、各ワーカー呼び出しで1メッセージあたりに行う処理量が増え、通信オーバーヘッドと計算量の比率を下げられます。

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個あります(10msで終わるものもあれば、500msかかるものもあります)。合計の経過時間を最小にするには、どのスケジューリング戦略を選ぶべきでしょうか。

まとめ: 並列コードのデバッグ

重要なポイント:

  • 並列化する前に関数を逐次実行してテストします。まず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時間対応のAIチューター)、R Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 R Academyコースには全4レッスンが含まれています。

「並列コードのデバッグと負荷分散」で何を学びますか?

並列ワーカーで発生するエラーを処理し、不均衡なワークロードを効果的に分散します。 ブラウザで直接実行するハンズオンコードでR Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

R Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのR Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。

「並列コードのデバッグと負荷分散」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このR Academyレッスンでコードを書いて実行できますか?

はい。すべてのR Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. parallel パッケージと detectCores()
  2. future フレームワーク
  3. furrr:purrr 操作の並列化
  4. 並列コードのデバッグと負荷分散
← R Academyに戻る