การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน
จัดการข้อผิดพลาดในตัวทำงานแบบขนานและปรับสมดุลภาระงานที่ไม่เท่ากันอย่างมีประสิทธิภาพ
การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน เป็นบทเรียน R Academy ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 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')คำถามที่พบบ่อย
บทเรียน “การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน” ฟรีหรือไม่
ใช่ — ข้อความเต็มของ “การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส R Academy ให้อัปเกรดเป็น CoddyKit PRO คอร์ส R Academy มีบทเรียนทั้งหมด 4 บทเรียน
คุณจะเรียนรู้อะไรในบทเรียน “การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน”
จัดการข้อผิดพลาดในตัวทำงานแบบขนานและปรับสมดุลภาระงานที่ไม่เท่ากันอย่างมีประสิทธิภาพ คุณปฏิบัติ R Academy ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน
คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน R Academy หรือไม่
ไม่จำเป็นต้องมีประสบการณ์มาก่อน R Academy บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน
บทเรียน “การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน” ใช้เวลานานแค่ไหน
บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย
ฉันเขียนและรันโค้ดในบทเรียน R Academy นี้ได้ไหม
ได้ บทเรียน R Academy ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ
บทเรียนทั้งหมดในหลักสูตรนี้
- แพ็กเกจ parallel และ detectCores()
- เฟรมเวิร์ก future
- furrr: การดำเนินการ purrr แบบขนาน
- การดีบักและการปรับสมดุลภาระงานของโค้ดแบบขนาน