【问题标题】:sys.sleep in apply and parSapplysys.sleep 在 apply 和 parSapply
【发布时间】:2020-06-20 15:09:35
【问题描述】:

我使用的代码会在给定时间内中断命令的执行,以免达到(crunchbase)API-Limit。

    managed_call <- function(f, events = 44L, every = 60L) {
    force(f)
    minute_ <- rep(NA, events)
    function(...) {       
    m_dif <- as.numeric(Sys.time() - minute_, units = "secs")
    minute_[!is.na(m_dif) & m_dif > every] <<- NA
    calls_remaining <- sum(is.na(minute_))
    if (!calls_remaining) {
    message("Close to API limit, pausing for ", 
    round(every - max(m_dif), 3), " seconds")
    Sys.sleep(every - max(m_dif))
    minute_[which.max(m_dif)] <- NA
    minute_[Position(is.na, minute_)] <<- Sys.time()
    f(...)
    } else {
    minute_[Position(is.na, minute_)] <<- Sys.time()
    f(...)
      }
     }
    }

当应用正常的 apply 或 lapply 命令时,这种代码和平会给我以下警告:

      Updated <- function(x){is.null(crunchbase_GET(x))}
  
      > abc <- unlist(lapply(websites,Updated))
      Close to API limit, pausing for 1.25 seconds
      Close to API limit, pausing for 1.119 seconds
      ...

但是,我尝试了使用 makeCluster 和 parSapply 的另一个选项:

library("parallel")

abc<- logical(100) 

Updated <- function(x){is.null(crunchbase_GET(x))}
cl <- makeCluster(detectCores(), type = "PSOCK")
clusterExport(cl, varlist = "websites")
clusterEvalQ(cl = cl, library(rcrunchbase))

abc <- parSapply(cl = cl, X = websites, FUN = Updated, USE.NAMES = FALSE)

警告消息现在不会出现。因此,我想知道是否执行了实际的 Sys.sleep() 命令,如果没有,是否有可能让我的代码使用 parSapply 运行。

非常抱歉,我无法提供可以针对这种特定情况复制的良好示例,因为需要 user_key 才能使用 rCrunchbase 并因此检索有关 API 限制等信息。

【问题讨论】:

    标签: r parallel-processing apply


    【解决方案1】:

    message 不会“转义”parSapply,它会丢失,catwarning 也是如此。将基本信息从cl 子进程传递给父进程的能力很困难。

    另一种选择(实际上是扩展,因为它们依赖于parallel)是futurefuture.apply,因为它们确实处理控制台输出。

    cl <- parallel::makeCluster(3)
    parallel::parLapply(cl, 1:3, function(i) { message("Hello: ", i+100); Sys.getpid(); })
    # [[1]]
    # [1] 22680
    # [[2]]
    # [1] 14504
    # [[3]]
    # [1] 27084
    

    但是future:

    library(future)        # plan, cluster
    library(future.apply)  # future_lapply
    # using the same 'cl'
    plan(cluster, workers = cl)
    future_lapply(1:3, function(i) { message("Hello: ", i+100); Sys.getpid(); })
    # Hello: 101
    # Hello: 102
    # Hello: 103
    # [[1]]
    # [1] 22680
    # [[2]]
    # [1] 14504
    # [[3]]
    # [1] 27084
    

    (变体可以证明catwarning 也可以转义子进程。)

    【讨论】:

    • 非常感谢!此解决方案是否也与 sys.Sleep() 兼容?
    • 我不明白为什么它不会,你试过发现问题了吗?
    • 我试过了,首先,我真的要说声谢谢,因为我真的很喜欢这种方法。但是,我应用标准 lapply 代码和 future_lapply 得到的结果完全不同,导致更多网站无法使用 future_lapply。此外,虽然 lapply 代码在处理列表中的前 1000 个 URL 时会停止很多次,但 future_lapply 命令似乎不会停止一次。我现在注意到错误应该出在 sys.sleep() 之前的一些步骤中。你能想象什么会在那里造成麻烦吗?
    • 使用的函数:Updated &lt;- function(x){is.null(crunchbase_GET(x))},未来的方法:cl &lt;- makeCluster(detectCores(), type = "PSOCK") plan(cluster, workers = cl) def &lt;- unlist(future_lapply(websites, Updated)) stopCluster(cl),lapply 的方法(我认为给出正确的结果):abc &lt;- unlist(lapply(websites,Updated))
    • 最后一条附加评论:我收到有关 HTTP 429 错误的各种警告,表明我确实通过了 API 限制
    猜你喜欢
    • 1970-01-01
    • 2021-11-02
    • 2011-09-05
    • 1970-01-01
    • 2015-06-18
    • 2014-12-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多