並列コードのデバッグと負荷分散
並列ワーカーで発生するエラーを処理し、不均衡なワークロードを効果的に分散します。
「並列コードのデバッグと負荷分散」は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フィードバックを取得できます。ローカル設定は不要です。