使用doParallel与cmdstanr时实时日志无法即时写入的问题
问题:并行贝叶斯A/B测试功效分析中实时日志写入失效
我使用cmdstanr实现了一个贝叶斯A/B测试的功效分析(样本量计算器),由于任务计算密集,尝试在foreach循环的每次迭代结束后用fwrite保存输出日志,但发现日志文件要等所有循环完成后才会写入。如果中途电脑崩溃,所有中间数据都会丢失。相同的日志写法在其他场景下正常工作,怀疑问题和cmdstanr的使用有关,求解决思路。
代码实现
library(cmdstanr) library(doParallel) library(data.table) #-------------------------- Parameters mde <- .007 # Minimum detectable effect (0.07%) mean_prior <- .01411629 # Past experiment mean (1.4%) sd_prior <- .103165 # Past experiment SD (10%) sample_size <- seq(1e6, 40e6, 1e6) |> rep(each = 1e2) # Range of sample sizes to test #-------------------------- Stan Model stan_model_code <- "data { int<lower=0> N; // Sample size vector[N] y; // Click-through rates real mu_prior; real sigma_prior; } parameters { real mu; // Mean click-through rate } model { mu ~ normal(mu_prior, sigma_prior); y ~ normal(mu, 0.1); // Assume 10% SD for individual user CTR } " # Compile the model stan_file <- write_stan_file(stan_model_code) stan_model <- cmdstan_model(stan_file) #-------------------------- Power Analysis Simulation Function sim <- \(n) { # Generate synthetic data for A and B groups y_A <- rnorm(n, mean = mean_prior, sd = 0.1) y_B <- rnorm(n, mean = mean_prior * (1 + mde), sd = 0.1) # Run Bayesian model for each group fit_A <- stan_model$sample(data = list(N = n, y = y_A, mu_prior = mean_prior, sigma_prior = sd_prior) , iter_sampling = 2000 , chains = 4 , parallel_chains = 1 , refresh = 0 ) fit_B <- stan_model$sample(data = list(N = n, y = y_B, mu_prior = mean_prior, sigma_prior = sd_prior) , iter_sampling = 2000 , chains = 4 , parallel_chains = 1 , refresh = 0 ) # Compute probability B > A return( list(n = n , power = mean(fit_B$draws("mu") > fit_A$draws("mu")) ) ) } #-------------------- set up cluster n_cores <- parallel::detectCores() - 1 cl <- makeCluster(n_cores) registerDoParallel(cl, cores = n_cores) #-------------------- simulation - parallel running h <- foreach(i = sample_size # , .combine = c , .packages = c("data.table", "cmdstanr") , .inorder = FALSE ) %dopar% { x <- sim(i) fwrite(x = data.table(time = Sys.time() , sample_size = i , power = x[[2]] ) , file = 'log - Bayesian cmdstanr.txt' , row.names = F , append = T ) return(x) }
解决思路
强制刷新文件缓冲区:
fwrite默认可能依赖系统缓冲,添加flush=TRUE参数可以强制将数据立即写入磁盘,修改后的代码片段:fwrite(x = data.table(time = Sys.time(), sample_size = i, power = x[[2]]), file = 'log - Bayesian cmdstanr.txt', row.names = F, append = T, flush = TRUE)也可以在
fwrite后调用flush.console()确保缓冲区清空。避免多进程文件写入冲突:并行子进程同时写入同一个文件可能触发系统的文件锁或缓冲延迟,改用文件连接的方式写入,每次打开、写入后立即关闭:
log_data <- paste(Sys.time(), i, x[[2]], sep = "\t") con <- file('log - Bayesian cmdstanr.txt', open = "a") writeLines(log_data, con) close(con)排查嵌套并行的资源冲突:你在
stan_model$sample中设置了parallel_chains = 1,外层又使用doParallel集群,可能存在嵌套并行导致的资源竞争或缓冲异常。可以尝试将parallel_chains设为0(禁用Stan内部并行),或者减少外层集群的核心数,避免资源过载。使用独立临时日志文件:为每个子进程生成独立的日志文件,最后再合并,避免多进程写入同一文件的冲突:
log_file <- paste0('log_', Sys.getpid(), '.txt') fwrite(x = data.table(time = Sys.time(), sample_size = i, power = x[[2]]), file = log_file, row.names = F)所有任务完成后,用
list.files找到所有临时日志,再合并成最终文件。
内容的提问来源于stack exchange,提问作者Sweepy Dodo
相关产品推荐
相关产品推荐

