You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 22:02:15